From c5663aa18b22e84d10c9d672514cf34b69dec2da Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 14:51:13 +0200 Subject: [PATCH] Let platform manifests go only on a confirmed collect, and copy again what a sweep took An unrecorded index or a copy in progress can name a platform the records do not see, so only a person's collect, after its dry run, takes an index's platforms, and only once every kept index of each repository it touches was read. A copy missing a platform is copied again, and a copy a build holds again is no longer recorded as collected (review of #144). --- cmd/mesh-controller/collect.go | 28 +++-- cmd/mesh-controller/mirrors.go | 24 ++++- cmd/mesh-controller/mirrors_test.go | 142 +++++++++++++++++++++++++ cmd/mesh-controller/registry_verbs.go | 2 +- internal/artifacts/platforms_test.go | 40 ++++++- internal/artifacts/store.go | 66 +++++++----- internal/builder/mirror.go | 44 +++++++- internal/builder/mirror_shared_test.go | 132 ++++++++++++++++++++++- internal/builder/mirror_test.go | 7 ++ internal/catalogue/verbs.go | 7 +- internal/inventory/builds.go | 11 ++ internal/inventory/mirrored_test.go | 30 ++++++ 12 files changed, 485 insertions(+), 48 deletions(-) diff --git a/cmd/mesh-controller/collect.go b/cmd/mesh-controller/collect.go index 861e28bf..936f02b7 100644 --- a/cmd/mesh-controller/collect.go +++ b/cmd/mesh-controller/collect.go @@ -31,6 +31,11 @@ type sweepBounds struct { most int // budget is the longest it keeps its caller waiting. budget time.Duration + // platforms is whether an index is let go of with its own platform manifests (novox/hq ADR 0257). + // Only a sweep a person confirmed, after its dry run listed what stays and names a platform: the + // records cannot see an unrecorded index, or a copy in progress, naming the same platform, and the + // store's tools can. + platforms bool } // afterBuild are the bounds of the sweep inside somebody's build. @@ -76,18 +81,27 @@ func sweep(ctx context.Context, inv *inventory.Inventory, store artifacts.Store, // as builds come, and is two HEADs each once done. A kept archive that could not be held stops // the sweep before it deletes anything: "everything kept is held" is the precondition the // collector's safety rests on, and a store refusing a hold would refuse the deletes too. - // **An index goes with its platform manifests, and a kept index keeps its own** (novox/hq ADR - // 0257): which of a repository's platform manifests a kept index names is read from the store - // before anything there is let go of. A keep set that cannot be read stops the sweep before it - // deletes anything, as a kept archive that cannot be held does. - if store.Spare == nil && inv != nil { + // **An index goes with its platform manifests only when a person confirmed it** (novox/hq ADR + // 0257), and a kept index keeps its own: which platform manifests the kept indexes name is read + // from the store for every repository the sweep will touch, before the first delete. A single + // read that fails stops the sweep before it deletes anything, as a kept archive that cannot be + // held does. The sweep after a build lets an index go alone, as it always has. + store.Spare = nil + if bounds.platforms { + if inv == nil { + r.Stopped = "no records to read which platform manifests a kept index names, so nothing was let go" + r.Left = len(references) + return r + } images, err := inv.KeptImages(ctx) + if err == nil { + store.Spare, err = store.SpareKeptIndexes(within, images, references) + } if err != nil { - r.Stopped = fmt.Sprintf("the images the mesh keeps could not be read, so nothing was let go: %v", err) + r.Stopped = fmt.Sprintf("which platform manifests a kept index names could not be read, so nothing was let go: %v", err) r.Left = len(references) return r } - store.Spare = store.SpareKeptIndexes(images) } wrote, missing, err := holdKept(within, store, kept) diff --git a/cmd/mesh-controller/mirrors.go b/cmd/mesh-controller/mirrors.go index 8b545715..901d3a8c 100644 --- a/cmd/mesh-controller/mirrors.go +++ b/cmd/mesh-controller/mirrors.go @@ -64,7 +64,9 @@ type mirrorsAnswer struct { } `json:"counts"` // Recorded, AlreadyRecorded and Refused answer a record: what was (or, a dry run, would be) // recorded, what the records already held, and what was refused, with why. - Recorded []string `json:"recorded,omitempty"` + Recorded []string `json:"recorded,omitempty"` + // Then says what recording does: a recorded copy is eligible, and the next sweep lets it go. + Then string `json:"then,omitempty"` AlreadyRecorded []string `json:"already_recorded,omitempty"` Refused map[string]string `json:"refused,omitempty"` } @@ -132,6 +134,15 @@ func recordMirrors(ctx context.Context, inv *inventory.Inventory, store artifact record = append(record, reference) } a.Recorded = record + if len(record) > 0 { + verb := "would become" + if real { + verb = "became" + } + a.Then = fmt.Sprintf("%d cop%s %s eligible: the next sweep (after any build, or collect) lets each go "+ + "unless one of the five kept builds of a module the mesh holds stood on it", len(record), + map[bool]string{true: "y", false: "ies"}[len(record) == 1], verb) + } if real && len(record) > 0 { if err := inv.RecordMirrored(ctx, "", why, record); err != nil { return a, fmt.Errorf("recording %d copies: %w", len(record), err) @@ -188,14 +199,16 @@ func mirrorsCommand(ctx context.Context, args []string) error { if err != nil { return err } - if *confirm { - f.record(ctx, "mirrors", []string{"--record", *record}) - } answer, err = recordMirrors(ctx, inv, artifacts.Store{Address: address}, splitReferences(*record), known, strings.TrimSpace(*f.why), *confirm) if err != nil { return err } + // The hand act is logged once the records say what it did, never before: an act that + // failed half-way is not logged as done. + if *confirm && len(answer.Recorded) > 0 { + f.record(ctx, "mirrors", append([]string{"--record"}, answer.Recorded...)) + } } if *asJSON { encoder := json.NewEncoder(os.Stdout) @@ -205,6 +218,9 @@ func mirrorsCommand(ctx context.Context, args []string) error { if answer.DryRun && len(answer.Recorded) > 0 { fmt.Println("a dry run: nothing was recorded (--confirm --why to record)") } + if answer.Then != "" { + fmt.Println(answer.Then) + } for _, r := range answer.Recorded { fmt.Printf(" record %s\n", r) } diff --git a/cmd/mesh-controller/mirrors_test.go b/cmd/mesh-controller/mirrors_test.go index 75da57c3..ae5e0c06 100644 --- a/cmd/mesh-controller/mirrors_test.go +++ b/cmd/mesh-controller/mirrors_test.go @@ -6,7 +6,9 @@ import ( "net/http/httptest" "slices" "strings" + "sync" "testing" + "time" "github.com/novox/mesh-controller/internal/artifacts" "github.com/novox/mesh-controller/internal/catalogue" @@ -111,6 +113,9 @@ func TestRecordingACopyTakesWhatTheStoreHoldsAndDeletesNothing(t *testing.T) { 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") } @@ -136,3 +141,140 @@ func TestRecordingACopyTakesWhatTheStoreHoldsAndDeletesNothing(t *testing.T) { } } } + +// 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 lets an index go alone: its platforms wait for a person's collect. +func TestTheSweepAfterABuildLetsAnIndexGoAlone(t *testing.T) { + inv := inventory.ForTest(t) + _, eligible, reg := twoCopies(t, inv) + store := reg.serve(t) + r := sweep(t.Context(), inv, store, []string{eligible}, nil, afterBuild) + if !slices.Equal(r.LetGo, []string{eligible}) || !slices.Equal(reg.deleted, []string{copyDigest(2)}) { + t.Fatalf("after a build: let go %v, deleted %v; want the index alone", r.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) + } +} diff --git a/cmd/mesh-controller/registry_verbs.go b/cmd/mesh-controller/registry_verbs.go index ab4dd994..1399944b 100644 --- a/cmd/mesh-controller/registry_verbs.go +++ b/cmd/mesh-controller/registry_verbs.go @@ -209,7 +209,7 @@ func collectCommand(ctx context.Context, args []string) error { if *most < 1 { return fmt.Errorf("collect: --most is at least 1, not %d", *most) } - bounds := sweepBounds{most: min(*most, collectMostAtMost), budget: collectBudget} + bounds := sweepBounds{most: min(*most, collectMostAtMost), budget: collectBudget, platforms: true} open, err := openStores(ctx) if err != nil { diff --git a/internal/artifacts/platforms_test.go b/internal/artifacts/platforms_test.go index 9f0793b2..e1784442 100644 --- a/internal/artifacts/platforms_test.go +++ b/internal/artifacts/platforms_test.go @@ -87,7 +87,11 @@ func TestAnIndexGoesWithItsPlatformsAndAKeptIndexKeepsItsOwn(t *testing.T) { repository := "upstream/docker.io/library/golang" old := catalogue.ArtifactStoreScheme + repository + "@" + digestN(1) kept := catalogue.ArtifactStoreScheme + repository + "@" + digestN(2) - store.Spare = store.SpareKeptIndexes([]string{kept}) + spare, err := store.SpareKeptIndexes(context.Background(), []string{kept}, []string{old}) + if err != nil { + t.Fatal(err) + } + store.Spare = spare if err := store.LetGo(context.Background(), old); err != nil { t.Fatal(err) @@ -113,7 +117,11 @@ func TestWithoutSparingAnIndexGoesAloneAndAnImageNamesNoPlatforms(t *testing.T) } // An image is one delete, with a Spare or without. s.deleted = nil - store.Spare = store.SpareKeptIndexes(nil) + spare, err := store.SpareKeptIndexes(context.Background(), nil, []string{catalogue.ArtifactStoreScheme + "web/app@" + digestN(5)}) + if err != nil { + t.Fatal(err) + } + store.Spare = spare if err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/app@"+digestN(5)); err != nil { t.Fatal(err) } @@ -122,7 +130,7 @@ func TestWithoutSparingAnIndexGoesAloneAndAnImageNamesNoPlatforms(t *testing.T) } // An index the store no longer holds is Gone, and nothing is deleted on its account. s.deleted = nil - err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/base@"+digestN(9)) + err = store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/base@"+digestN(9)) if !errors.Is(err, Gone) || len(s.deleted) != 0 { t.Fatalf("an index the store does not hold: %v, deleted %v", err, s.deleted) } @@ -140,3 +148,29 @@ func TestASpareThatCannotBeReadKeepsTheIndex(t *testing.T) { t.Fatalf("deleted %v without knowing what a kept index names", s.deleted) } } + +// The spare list is read whole before the sweep: a kept index that cannot be read is an error before +// anything is deleted, and a repository not read before keeps its index. +func TestASpareListIsReadWholeBeforeAnythingGoes(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + t.Cleanup(server.Close) + store := Store{Address: strings.TrimPrefix(server.URL, "http://")} + kept := catalogue.ArtifactStoreScheme + "web/base@" + digestN(2) + old := catalogue.ArtifactStoreScheme + "web/base@" + digestN(1) + if _, err := store.SpareKeptIndexes(context.Background(), []string{kept}, []string{old}); err == nil { + t.Fatal("a kept index the store would not answer for was taken as naming nothing") + } + + s := &indexStore{indexes: map[string][]string{digestN(1): {digestN(11)}}, images: map[string]bool{digestN(11): true}} + good := s.serve(t) + spare, err := good.SpareKeptIndexes(context.Background(), nil, []string{catalogue.ArtifactStoreScheme + "other/repo@" + digestN(3)}) + if err != nil { + t.Fatal(err) + } + good.Spare = spare + if err := good.LetGo(context.Background(), old); err == nil || len(s.deleted) != 0 { + t.Fatalf("an index in a repository not read before the sweep was let go of: %v, deleted %v", err, s.deleted) + } +} diff --git a/internal/artifacts/store.go b/internal/artifacts/store.go index 0b98ed7c..253c428c 100644 --- a/internal/artifacts/store.go +++ b/internal/artifacts/store.go @@ -169,44 +169,54 @@ func (s Store) Platforms(ctx context.Context, repository, digest string) ([]stri return out, nil } -// SpareKeptIndexes is a Spare for a sweep: the platform manifests the kept references of each -// repository name, read from the store once per repository and remembered for the sweep. A kept -// reference itself is spared too. A kept index the store does not hold names nothing. -func (s Store) SpareKeptIndexes(kept []string) func(ctx context.Context, repository string) (map[string]bool, error) { - byRepository := map[string][]string{} - for _, reference := range kept { +// SpareKeptIndexes is a Spare for a sweep, read whole before it starts: for every repository an +// image reference of `touching` is in, each kept manifest there and every platform manifest a kept +// index there names. Any read that fails is an error, and the sweep lets nothing go. A repository not +// read here is answered with an error, so an index in it is kept. A kept index the store no longer +// holds names nothing. +func (s Store) SpareKeptIndexes(ctx context.Context, kept, touching []string) (func(context.Context, string) (map[string]bool, error), error) { + imageIn := func(reference string) (repository, digest string, ok bool) { path, ours := catalogue.InArtifactStore(reference) if !ours { - continue + return "", "", false } repository, kind, digest, err := split(path) if err != nil || kind != "manifests" { - continue + return "", "", false } - byRepository[repository] = append(byRepository[repository], digest) + return repository, digest, true } read := map[string]map[string]bool{} - return func(ctx context.Context, repository string) (map[string]bool, error) { - if spared, done := read[repository]; done { - return spared, nil + for _, reference := range touching { + if repository, _, ok := imageIn(reference); ok { + read[repository] = map[string]bool{} } - spared := map[string]bool{} - for _, digest := range byRepository[repository] { - spared[digest] = true - children, err := s.Platforms(ctx, repository, digest) - if errors.Is(err, Gone) { - continue - } - if err != nil { - return nil, err - } - for _, c := range children { - spared[c] = true - } - } - read[repository] = spared - return spared, nil } + for _, reference := range kept { + repository, digest, ok := imageIn(reference) + spared, touched := read[repository] + if !ok || !touched { + continue + } + spared[digest] = true + children, err := s.Platforms(ctx, repository, digest) + if errors.Is(err, Gone) { + continue + } + if err != nil { + return nil, fmt.Errorf("reading the kept %s@%s: %w", repository, digest, err) + } + for _, c := range children { + spared[c] = true + } + } + return func(_ context.Context, repository string) (map[string]bool, error) { + spared, done := read[repository] + if !done { + return nil, fmt.Errorf("%s was not read before the sweep began", repository) + } + return spared, nil + }, nil } // remove asks the store to delete what is at url. Gone when it has no such thing. diff --git a/internal/builder/mirror.go b/internal/builder/mirror.go index e449e3ce..8840e3d5 100644 --- a/internal/builder/mirror.go +++ b/internal/builder/mirror.go @@ -232,7 +232,11 @@ func (r Registry) MirrorBase(ctx context.Context, from, repository, former strin // on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base // lives there failed on a copy it did not need. if strings.HasPrefix(where.reference, "sha256:") { - held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference, manifestAccept) + // **Held whole, not only its index** (novox/hq ADR 0257): every platform manifest the index + // names is asked about too. A sweep that stopped half-way, or that ran while this build was + // copying, may have left an index naming a platform the store no longer holds; such a copy + // is made again, and the copy puts back what is missing. + held, err := r.holdsWhole(ctx, repository, where.reference) if err != nil { return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err) } @@ -242,7 +246,7 @@ func (r Registry) MirrorBase(ctx context.Context, from, repository, former strin } src := &source{client: r.client()} if former != "" && former != repository && strings.HasPrefix(where.reference, "sha256:") { - held, err := r.has(ctx, "http://"+r.Address+"/v2/"+former+"/manifests/"+where.reference, manifestAccept) + held, err := r.holdsWhole(ctx, former, where.reference) if err != nil { return "", fmt.Errorf("asking %s whether %s holds %s: %w", r.Address, former, from, err) } @@ -259,6 +263,42 @@ func (r Registry) MirrorBase(ctx context.Context, from, repository, former strin return r.Address + "/" + repository + "@" + digest, nil } +// holdsWhole is whether this registry holds the manifest under the repository and, for an index, every +// manifest it names. +func (r Registry) holdsWhole(ctx context.Context, repository, digest string) (bool, error) { + url := "http://" + r.Address + "/v2/" + repository + "/manifests/" + held, err := r.has(ctx, url+digest, manifestAccept) + if err != nil || !held { + return false, err + } + request, err := http.NewRequestWithContext(ctx, http.MethodGet, url+digest, nil) + if err != nil { + return false, err + } + request.Header.Set("Accept", manifestAccept) + response, err := r.client().Do(request) + if err != nil { + return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err) + } + defer response.Body.Close() + if response.StatusCode != http.StatusOK { + return false, nil + } + var document struct { + Manifests []descriptor `json:"manifests"` + } + if err := json.NewDecoder(io.LimitReader(response.Body, 4<<20)).Decode(&document); err != nil { + return false, fmt.Errorf("%s@%s is not a manifest: %w", repository, digest, err) + } + for _, m := range document.Manifests { + held, err := r.has(ctx, url+m.Digest, manifestAccept) + if err != nil || !held { + return false, err + } + } + return true, nil +} + // copyManifest copies one manifest document and everything it names, and returns its digest. An // index is copied by copying each manifest it names first, so the index never points at something // the registry does not hold yet. diff --git a/internal/builder/mirror_shared_test.go b/internal/builder/mirror_shared_test.go index d7824d02..45d21695 100644 --- a/internal/builder/mirror_shared_test.go +++ b/internal/builder/mirror_shared_test.go @@ -43,6 +43,9 @@ type aRegistryOfRepositories struct { uploads int mounts int gets int + // noMount answers every mount with an ordinary upload's location, as a registry that cannot + // mount across repositories does. + noMount bool } func newRegistryOfRepositories() *aRegistryOfRepositories { @@ -102,8 +105,11 @@ func (m *aRegistryOfRepositories) handler() http.Handler { return } w.WriteHeader(http.StatusOK) + case kind == "manifests" && r.Method == http.MethodDelete: + delete(m.manifests, repository+"@"+rest) + w.WriteHeader(http.StatusAccepted) case kind == "blobs/uploads" && r.Method == http.MethodPost: - if digest, from := r.URL.Query().Get("mount"), r.URL.Query().Get("from"); digest != "" && m.links[from][digest] { + if digest, from := r.URL.Query().Get("mount"), r.URL.Query().Get("from"); digest != "" && m.links[from][digest] && !m.noMount { m.link(repository, digest) m.mounts++ w.WriteHeader(http.StatusCreated) @@ -243,3 +249,127 @@ func TestABuildSaysWhatItMirroredEvenWhenItFails(t *testing.T) { t.Fatalf("a build said it mirrored %v, want [%s]", ok.Mirrored, want) } } + +// held puts one image in a repository the old way and answers its index digest and the registry. +func heldUnderAModule(t *testing.T) (*aRegistryOfRepositories, Registry, string, string, map[string][]byte, func()) { + t.Helper() + src, indexDigest, srcBlobs := anUpstreamRegistry(t) + t.Cleanup(src.Close) + dst := newRegistryOfRepositories() + dstServer := httptest.NewServer(dst.handler()) + t.Cleanup(dstServer.Close) + r := Registry{Address: strings.TrimPrefix(dstServer.URL, "http://"), HTTP: src.Client()} + host := strings.TrimPrefix(src.URL, "http://") + if _, err := r.MirrorImage(context.Background(), host+"/library/thing:latest", "hello-web/on-thing_base"); err != nil { + t.Fatal(err) + } + return dst, r, host + "/library/thing@" + indexDigest, indexDigest, srcBlobs, src.Close +} + +// An index whose platform manifests are not all held is not "already held": it is copied again, and the +// copy puts back what is missing — what a sweep that stopped half-way, or ran while a build was copying, +// leaves behind. +func TestAnIndexMissingAPlatformIsCopiedAgain(t *testing.T) { + dst, r, from, indexDigest, _, _ := heldUnderAModule(t) + repository, _ := MirrorRepository(from) + if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil { + t.Fatal(err) + } + // A sweep lets one platform go, between this build's copy and the next. + var platform string + for key := range dst.manifests { + if strings.HasPrefix(key, repository+"@") && !strings.HasSuffix(key, indexDigest) { + platform = key + break + } + } + delete(dst.manifests, platform) + if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil { + t.Fatal(err) + } + if _, back := dst.manifests[platform]; !back { + t.Fatalf("an index missing %s was taken as held", platform) + } +} + +// The same, for the module's former copy: a former copy missing a platform is not a source; upstream +// is asked instead. +func TestAFormerCopyMissingAPlatformIsNotTheSource(t *testing.T) { + dst, r, from, indexDigest, _, _ := heldUnderAModule(t) + for key := range dst.manifests { + if strings.HasPrefix(key, "hello-web/on-thing_base@") && !strings.HasSuffix(key, indexDigest) { + delete(dst.manifests, key) + break + } + } + repository, _ := MirrorRepository(from) + if _, err := r.MirrorBase(context.Background(), from, repository, "hello-web/on-thing_base"); err != nil { + t.Fatal(err) + } + if dst.mounts != 0 { + t.Fatalf("a former copy missing a platform was mounted from (%d mounts)", dst.mounts) + } + if n := countPrefix(dst.manifests, repository+"@"); n != 3 { + t.Fatalf("%d manifests copied from upstream, want 3", n) + } +} + +// A registry that answers a mount with an upload's location (202) gets the blob moved as before. +func TestAMountRefusedFallsBackToAnUpload(t *testing.T) { + dst, r, from, _, srcBlobs, _ := heldUnderAModule(t) + dst.noMount = true + uploaded := dst.uploads + repository, _ := MirrorRepository(from) + reference, err := r.MirrorBase(context.Background(), from, repository, "hello-web/on-thing_base") + if err != nil { + t.Fatal(err) + } + if !strings.HasSuffix(reference, repository+"@"+strings.SplitN(from, "@", 2)[1]) { + t.Fatalf("pinned as %s", reference) + } + if dst.mounts != 0 || dst.uploads-uploaded != len(srcBlobs) { + t.Fatalf("%d mounts, %d uploads; want every blob uploaded from the former copy", dst.mounts, dst.uploads-uploaded) + } + for digest := range srcBlobs { + if !dst.links[repository][digest] { + t.Fatalf("blob %s is not linked into %s", digest, repository) + } + } +} + +// A sweep and a build at once: while builds copy the base, platforms are let go of under them. Each +// build that returns has the copy whole. +func TestABuildCopyingWhileASweepLetsGoEndsWhole(t *testing.T) { + dst, r, from, indexDigest, _, _ := heldUnderAModule(t) + repository, _ := MirrorRepository(from) + if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil { + t.Fatal(err) + } + done := make(chan struct{}) + go func() { + defer close(done) + for i := 0; i < 50; i++ { + dst.mu.Lock() + for key := range dst.manifests { + if strings.HasPrefix(key, repository+"@") && !strings.HasSuffix(key, indexDigest) { + delete(dst.manifests, key) + break + } + } + dst.mu.Unlock() + } + }() + for i := 0; i < 20; i++ { + if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil { + t.Fatal(err) + } + } + <-done + // The sweep is over; the next build finds what it left and makes the copy whole. + if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil { + t.Fatal(err) + } + if n := countPrefix(dst.manifests, repository+"@"); n != 3 { + t.Fatalf("after the sweep and the builds, %d manifests in %s, want 3", n, repository) + } +} diff --git a/internal/builder/mirror_test.go b/internal/builder/mirror_test.go index 9a17b187..03f4b519 100644 --- a/internal/builder/mirror_test.go +++ b/internal/builder/mirror_test.go @@ -117,6 +117,13 @@ func (m *theMeshsRegistry) handler() http.Handler { } else { w.WriteHeader(http.StatusNotFound) } + case r.Method == http.MethodGet && strings.Contains(r.URL.Path, "/manifests/"): + body, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]] + if !ok { + w.WriteHeader(http.StatusNotFound) + return + } + _, _ = w.Write(body) case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"): if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok { w.WriteHeader(http.StatusOK) diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index 8c228c1b..31de26ba 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -449,7 +449,9 @@ var ControllerVerbs = []Verb{ {Name: "collect", Description: "The artifact store's sweep, now: what the mesh made and keeps for no reason " + "is let go of. Without confirm a dry run — every kept archive asked about with HEADs, what would be let go " + "listed, nothing held or deleted. With confirm and why: every kept archive held first, then each eligible " + - "artifact let go of, oldest first, and recorded collected — a hand act, recorded with its why. Bounded by " + + "artifact let go of, oldest first, and recorded collected — a hand act, recorded with its why. A confirmed " + + "collect lets an index go with the platform manifests no kept index names (read first; any read that " + + "fails lets nothing go); the sweep after a build lets an index go alone (novox/hq ADR 0257). Bounded by " + "most (500 by default, at most 5000) and 45 seconds; what is left is said. Bytes are reclaimed by the " + "store's nightly collector. Never touches a digest the mesh did not record (novox/hq ADR 0189, ADR 0251).", Input: schema(map[string]string{ @@ -463,7 +465,8 @@ var ControllerVerbs = []Verb{ "before copies were recorded, or by a build that failed before then) to record as the mesh's, so the " + "sweep keeps or lets go of them like any other; only a repository the mirror writes is taken " + "(upstream//, or a module's /on-), and each must be in the store. Without " + - "confirm a dry run; with confirm and why, recorded — a hand act. Records only: nothing is deleted " + + "confirm a dry run; with confirm and why, recorded — a hand act. A recorded copy is eligible: the next " + + "sweep lets it go unless a kept build stood on it. Records only: this verb deletes nothing " + "(novox/hq ADR 0257).", Input: schema(map[string]string{ "record": "references to record, artifact-store://@sha256:, separated by spaces or commas", diff --git a/internal/inventory/builds.go b/internal/inventory/builds.go index 346a17e2..08a48235 100644 --- a/internal/inventory/builds.go +++ b/internal/inventory/builds.go @@ -131,6 +131,11 @@ func (i *Inventory) RecordBuild(ctx context.Context, b Build) error { // RecordMirrored keeps that the mesh copied these bases into the artifact store (novox/hq ADR 0257): // by the build that said so, or by a person who says why. The first record of a copy stands. +// +// **A copy a build says it holds is held again.** The builder answers a copy only once the store holds +// it whole, copying again what was let go of; so a copy the records say was collected, and a build now +// says it copied, is no longer collected. Left collected, it would stay in the store and out of every +// sweep for ever. func (i *Inventory) RecordMirrored(ctx context.Context, build, why string, references []string) error { var by *string if build != "" { @@ -145,6 +150,12 @@ func (i *Inventory) RecordMirrored(ctx context.Context, build, why string, refer on conflict (reference) do nothing`, reference, by, why); err != nil { return err } + if build != "" { + if _, err := i.store.Pool().Exec(ctx, + `delete from artifact_collected where reference = $1`, reference); err != nil { + return err + } + } } return nil } diff --git a/internal/inventory/mirrored_test.go b/internal/inventory/mirrored_test.go index dea3b398..3fb13464 100644 --- a/internal/inventory/mirrored_test.go +++ b/internal/inventory/mirrored_test.go @@ -179,3 +179,33 @@ func TestTheCopiesTheRecordsAlreadyNameAreRecorded(t *testing.T) { t.Fatalf("recorded %s by %s", reference, build) } } + +// A copy let go of and copied again by a later build is held again: its record is back, and the sweep +// decides about it as before. +func TestACopyLetGoOfAndCopiedAgainIsHeldAgain(t *testing.T) { + inv := fresh(t) + ctx := context.Background() + holding(t, inv, "proxy") + golang := mirrorRef("upstream/docker.io/library/golang", 1) + builtOn(t, inv, "b01", "proxy", 1, []string{golang}, golang) + if err := inv.MarkCollected(ctx, []string{golang}); err != nil { + t.Fatal(err) + } + if s := stateOf(t, inv, golang); s.State != ArtifactCollected { + t.Fatalf("a let-go copy reads %+v", s) + } + builtOn(t, inv, "b02", "proxy", 2, []string{golang}, golang) + if s := stateOf(t, inv, golang); s.State != ArtifactKept || !slices.Equal(s.Why, []string{KeptStoodOn}) { + t.Fatalf("a copy copied again reads %+v; want kept again", s) + } + // A person's record never undoes a collection: only a build that copied it says it is there. + if err := inv.MarkCollected(ctx, []string{golang}); err != nil { + t.Fatal(err) + } + if err := inv.RecordMirrored(ctx, "", "again", []string{golang}); err != nil { + t.Fatal(err) + } + if s := stateOf(t, inv, golang); s.State != ArtifactCollected { + t.Fatalf("a person's record brought a collected copy back: %+v", s) + } +}