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).
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user