From 412599de9a5bd848b415ff9635b0bda8639f9073 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 21 Sep 2026 20:42:30 +0200 Subject: [PATCH] An upstream image is copied between registries, never through a machine's image store A published image is an index over several architectures; pulled, the runtime's store keeps the index and refuses to push one platform out of it. The builder now reads the index and every manifest it names over the registry API, with the anonymous bearer token the public hub hands out, moves each blob by digest into the mesh's registry, puts the manifests and then the index under the module's repository, and pins the index. Genesis, with no registry to copy into, keeps the pull (novox/hq 04-ISSUES/046, ADR 0096). --- internal/builder/mirror.go | 318 ++++++++++++++++++++++++++++++++ internal/builder/mirror_test.go | 210 +++++++++++++++++++++ 2 files changed, 528 insertions(+) create mode 100644 internal/builder/mirror.go create mode 100644 internal/builder/mirror_test.go diff --git a/internal/builder/mirror.go b/internal/builder/mirror.go new file mode 100644 index 0000000..2f6b675 --- /dev/null +++ b/internal/builder/mirror.go @@ -0,0 +1,318 @@ +package builder + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" +) + +// An upstream image is copied between registries, never through a machine's image store +// (novox/hq 04-ISSUES/046, ADR 0096). +// +// A published image is ordinarily an index over several architectures. Pulling it leaves the index +// in the runtime's store, and pushing one platform out of that store is what the runtime refuses — +// every variant of pull-then-push was tried and failed the same way. A copy never needs a platform: +// it moves what is there. Read the index, read each manifest it names, put each blob by digest, +// put the manifests, put the index under the module's repository — the registry API is enough, and +// the runtime's image store is never involved. + +// Mirrorer copies an upstream image into the mesh's own registry, whole. +type Mirrorer interface { + MirrorImage(ctx context.Context, from, repository string) (string, error) +} + +const ( + mediaIndexOCI = "application/vnd.oci.image.index.v1+json" + mediaIndexDocker = "application/vnd.docker.distribution.manifest.list.v2+json" + mediaManifestOCI = "application/vnd.oci.image.manifest.v1+json" + mediaManifestDocker = "application/vnd.docker.distribution.manifest.v2+json" +) + +var manifestAccept = strings.Join([]string{mediaIndexOCI, mediaIndexDocker, mediaManifestOCI, mediaManifestDocker}, ", ") + +// upstream is where an image lives, as the registry API addresses it. +type upstream struct { + // base is the scheme and host, e.g. https://registry-1.docker.io. + base string + // repository is the path under /v2/, e.g. library/alpine. + repository string + // reference is a tag or a digest. + reference string +} + +// parseReference splits `[host/]repo[:tag][@digest]` the way a container runtime does: no host +// means the public hub, and a single-segment repository there lives under `library/`. +func parseReference(ref string) (upstream, error) { + name, reference := ref, "latest" + if at := strings.Index(ref, "@"); at >= 0 { + name, reference = ref[:at], ref[at+1:] + } else if colon := strings.LastIndex(ref, ":"); colon > strings.LastIndex(ref, "/") { + name, reference = ref[:colon], ref[colon+1:] + } + if name == "" || reference == "" { + return upstream{}, fmt.Errorf("%q is not an image reference", ref) + } + host, repository := "docker.io", name + if slash := strings.Index(name, "/"); slash >= 0 && strings.ContainsAny(name[:slash], ".:") { + host, repository = name[:slash], name[slash+1:] + } else if slash >= 0 && name[:slash] == "localhost" { + host, repository = name[:slash], name[slash+1:] + } + if host == "docker.io" { + host = "registry-1.docker.io" + if !strings.Contains(repository, "/") { + repository = "library/" + repository + } + } + scheme := "https://" + if strings.HasPrefix(host, "localhost") || strings.HasPrefix(host, "127.") { + scheme = "http://" + } + return upstream{base: scheme + host, repository: repository, reference: reference}, nil +} + +// source reads from one upstream registry, taking a bearer token where the registry asks for one. +type source struct { + client *http.Client + token string +} + +// get fetches a registry URL, answering a bearer challenge once with an anonymous token — which is +// how the public hub serves public images, and every registry the catalogue names does the same. +func (s *source) get(ctx context.Context, url, accept string) (*http.Response, error) { + for attempt := 0; attempt < 2; attempt++ { + request, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return nil, err + } + if accept != "" { + request.Header.Set("Accept", accept) + } + if s.token != "" { + request.Header.Set("Authorization", "Bearer "+s.token) + } + response, err := s.client.Do(request) + if err != nil { + return nil, err + } + if response.StatusCode != http.StatusUnauthorized || attempt == 1 { + return response, nil + } + challenge := response.Header.Get("WWW-Authenticate") + response.Body.Close() + token, err := s.tokenFor(ctx, challenge) + if err != nil { + return nil, err + } + s.token = token + } + return nil, fmt.Errorf("unreachable") +} + +// tokenFor answers `Bearer realm="…",service="…",scope="…"` with an anonymous token request. +func (s *source) tokenFor(ctx context.Context, challenge string) (string, error) { + if !strings.HasPrefix(challenge, "Bearer ") { + return "", fmt.Errorf("the registry asks for %q, and this copies public images anonymously", challenge) + } + fields := map[string]string{} + for _, part := range strings.Split(challenge[len("Bearer "):], ",") { + key, value, found := strings.Cut(strings.TrimSpace(part), "=") + if found { + fields[key] = strings.Trim(value, `"`) + } + } + realm := fields["realm"] + if realm == "" { + return "", fmt.Errorf("the registry's challenge names no realm: %q", challenge) + } + url := realm + "?service=" + fields["service"] + if scope := fields["scope"]; scope != "" { + url += "&scope=" + scope + } + request, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return "", err + } + response, err := s.client.Do(request) + if err != nil { + return "", fmt.Errorf("cannot get a token from %s: %w", realm, err) + } + defer response.Body.Close() + var issued struct { + Token string `json:"token"` + AccessToken string `json:"access_token"` + } + if err := json.NewDecoder(response.Body).Decode(&issued); err != nil { + return "", fmt.Errorf("%s answered with something that is not a token: %w", realm, err) + } + if issued.Token == "" { + issued.Token = issued.AccessToken + } + if issued.Token == "" { + return "", fmt.Errorf("%s issued no token", realm) + } + return issued.Token, nil +} + +// descriptor is what an index or a manifest names: a blob or another manifest, by digest. +type descriptor struct { + MediaType string `json:"mediaType"` + Digest string `json:"digest"` + Size int64 `json:"size"` +} + +// MirrorImage copies `from` — an index or a single manifest, by tag or digest — into this registry +// under `repository`, and returns the reference the mesh will pin: this registry, the repository, +// and the digest of the document that was put last, which is the index where there is one. +func (r Registry) MirrorImage(ctx context.Context, from, repository string) (string, error) { + where, err := parseReference(from) + if err != nil { + return "", err + } + src := &source{client: r.client()} + digest, err := r.copyManifest(ctx, src, where, where.reference, repository) + if err != nil { + return "", fmt.Errorf("copying %s into %s/%s: %w", from, r.Address, repository, err) + } + return r.Address + "/" + repository + "@" + digest, 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. +func (r Registry) copyManifest(ctx context.Context, src *source, where upstream, reference, repository string) (string, error) { + response, err := src.get(ctx, where.base+"/v2/"+where.repository+"/manifests/"+reference, manifestAccept) + if err != nil { + return "", err + } + defer response.Body.Close() + if response.StatusCode != http.StatusOK { + said, _ := io.ReadAll(io.LimitReader(response.Body, 2048)) + return "", fmt.Errorf("%s/%s@%s: %s %s", where.base, where.repository, reference, response.Status, strings.TrimSpace(string(said))) + } + body, err := io.ReadAll(response.Body) + if err != nil { + return "", err + } + mediaType := response.Header.Get("Content-Type") + if semi := strings.Index(mediaType, ";"); semi >= 0 { + mediaType = mediaType[:semi] + } + var document struct { + MediaType string `json:"mediaType"` + Manifests []descriptor `json:"manifests"` + Config *descriptor `json:"config"` + Layers []descriptor `json:"layers"` + } + if err := json.Unmarshal(body, &document); err != nil { + return "", fmt.Errorf("%s is not a manifest: %w", reference, err) + } + if mediaType == "" || mediaType == "application/json" { + mediaType = document.MediaType + } + + switch mediaType { + case mediaIndexOCI, mediaIndexDocker: + // The manifests first, each by its digest; the index that names them last. + for _, m := range document.Manifests { + if _, err := r.copyManifest(ctx, src, where, m.Digest, repository); err != nil { + return "", err + } + } + case mediaManifestOCI, mediaManifestDocker: + blobs := append([]descriptor{}, document.Layers...) + if document.Config != nil { + blobs = append(blobs, *document.Config) + } + for _, b := range blobs { + if err := r.copyBlob(ctx, src, where, b.Digest, repository); err != nil { + return "", err + } + } + default: + return "", fmt.Errorf("%s is a %q, which is neither an image index nor an image manifest", reference, mediaType) + } + + digest := "sha256:" + hexOf(sha256.Sum256(body)) + put, err := http.NewRequestWithContext(ctx, http.MethodPut, + "http://"+r.Address+"/v2/"+repository+"/manifests/"+digest, bytes.NewReader(body)) + if err != nil { + return "", err + } + put.Header.Set("Content-Type", mediaType) + done, err := r.client().Do(put) + if err != nil { + return "", fmt.Errorf("cannot put a manifest into %s: %w", r.Address, err) + } + defer done.Body.Close() + if done.StatusCode != http.StatusCreated { + said, _ := io.ReadAll(io.LimitReader(done.Body, 2048)) + return "", fmt.Errorf("%s refused the manifest %s: %s %s", r.Address, digest, done.Status, strings.TrimSpace(string(said))) + } + return digest, nil +} + +// copyBlob moves one blob by digest, unless the registry already holds it — blobs are immutable +// and content-named, so "already there" is the whole check. +func (r Registry) copyBlob(ctx context.Context, src *source, where upstream, digest, repository string) error { + base := "http://" + r.Address + "/v2/" + repository + if there, err := r.has(ctx, base+"/blobs/"+digest); err != nil { + return err + } else if there { + return nil + } + response, err := src.get(ctx, where.base+"/v2/"+where.repository+"/blobs/"+digest, "") + if err != nil { + return err + } + defer response.Body.Close() + if response.StatusCode != http.StatusOK { + return fmt.Errorf("%s/%s: blob %s: %s", where.base, where.repository, digest, response.Status) + } + + start, err := http.NewRequestWithContext(ctx, http.MethodPost, base+"/blobs/uploads/", nil) + if err != nil { + return err + } + begun, err := r.client().Do(start) + if err != nil { + return fmt.Errorf("cannot start an upload to %s: %w", base, err) + } + begun.Body.Close() + if begun.StatusCode != http.StatusAccepted { + return fmt.Errorf("%s answered %s when asked where to put a blob", base, begun.Status) + } + location := begun.Header.Get("Location") + if location == "" { + return fmt.Errorf("%s accepted an upload and said nowhere to put it", base) + } + if strings.HasPrefix(location, "/") { + location = "http://" + r.Address + location + } + put, err := http.NewRequestWithContext(ctx, http.MethodPut, location+separator(location)+"digest="+digest, response.Body) + if err != nil { + return err + } + put.Header.Set("Content-Type", "application/octet-stream") + if response.ContentLength > 0 { + put.ContentLength = response.ContentLength + } + done, err := r.client().Do(put) + if err != nil { + return fmt.Errorf("cannot upload blob %s: %w", digest, err) + } + defer done.Body.Close() + if done.StatusCode != http.StatusCreated { + said, _ := io.ReadAll(io.LimitReader(done.Body, 2048)) + return fmt.Errorf("%s refused blob %s: %s %s", base, digest, done.Status, strings.TrimSpace(string(said))) + } + return nil +} + +func hexOf(sum [32]byte) string { return hex.EncodeToString(sum[:]) } diff --git a/internal/builder/mirror_test.go b/internal/builder/mirror_test.go new file mode 100644 index 0000000..38fe594 --- /dev/null +++ b/internal/builder/mirror_test.go @@ -0,0 +1,210 @@ +package builder + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" +) + +// An upstream image is copied between registries, never through a machine's image store +// (novox/hq 04-ISSUES/046, ADR 0096): the index, every manifest it names, every blob — moved by +// digest, and the index put last under the module's repository. + +func digestOf(b []byte) string { + sum := sha256.Sum256(b) + return "sha256:" + hex.EncodeToString(sum[:]) +} + +// anUpstreamRegistry serves one image as an index over two platforms, behind an anonymous bearer +// challenge the way the public hub does, and records what was fetched. +func anUpstreamRegistry(t *testing.T) (*httptest.Server, string, map[string][]byte) { + t.Helper() + blobs := map[string][]byte{} + manifests := map[string][]byte{} + put := func(kind, mediaType string, layer []byte) string { + config := []byte(`{"architecture":"` + kind + `"}`) + blobs[digestOf(config)] = config + blobs[digestOf(layer)] = layer + m, _ := json.Marshal(map[string]any{ + "schemaVersion": 2, "mediaType": mediaType, + "config": map[string]any{"mediaType": "application/vnd.oci.image.config.v1+json", "digest": digestOf(config), "size": len(config)}, + "layers": []map[string]any{{"mediaType": "application/vnd.oci.image.layer.v1.tar+gzip", "digest": digestOf(layer), "size": len(layer)}}, + }) + manifests[digestOf(m)] = m + return digestOf(m) + } + amd := put("amd64", mediaManifestOCI, []byte("amd64 layer bytes")) + arm := put("arm64", mediaManifestOCI, []byte("arm64 layer bytes")) + index, _ := json.Marshal(map[string]any{ + "schemaVersion": 2, "mediaType": mediaIndexOCI, + "manifests": []map[string]any{ + {"mediaType": mediaManifestOCI, "digest": amd, "size": len(manifests[amd]), "platform": map[string]string{"os": "linux", "architecture": "amd64"}}, + {"mediaType": mediaManifestOCI, "digest": arm, "size": len(manifests[arm]), "platform": map[string]string{"os": "linux", "architecture": "arm64"}}, + }, + }) + manifests["latest"] = index + manifests[digestOf(index)] = index + + var server *httptest.Server + server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/token" { + _, _ = w.Write([]byte(`{"token":"anonymous-token"}`)) + return + } + if r.Header.Get("Authorization") != "Bearer anonymous-token" { + w.Header().Set("WWW-Authenticate", `Bearer realm="`+server.URL+`/token",service="test",scope="repository:library/thing:pull"`) + w.WriteHeader(http.StatusUnauthorized) + return + } + switch { + case strings.HasPrefix(r.URL.Path, "/v2/library/thing/manifests/"): + ref := strings.TrimPrefix(r.URL.Path, "/v2/library/thing/manifests/") + body, ok := manifests[ref] + if !ok { + w.WriteHeader(http.StatusNotFound) + return + } + var typed struct { + MediaType string `json:"mediaType"` + } + _ = json.Unmarshal(body, &typed) + w.Header().Set("Content-Type", typed.MediaType) + _, _ = w.Write(body) + case strings.HasPrefix(r.URL.Path, "/v2/library/thing/blobs/"): + body, ok := blobs[strings.TrimPrefix(r.URL.Path, "/v2/library/thing/blobs/")] + if !ok { + w.WriteHeader(http.StatusNotFound) + return + } + _, _ = w.Write(body) + default: + w.WriteHeader(http.StatusNotFound) + } + })) + return server, digestOf(index), blobs +} + +// theMeshsRegistry accepts blobs and manifests the way a registry does, and remembers them. +type theMeshsRegistry struct { + mu sync.Mutex + blobs map[string][]byte + manifests map[string][]byte + uploads int +} + +func (m *theMeshsRegistry) handler() http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + m.mu.Lock() + defer m.mu.Unlock() + switch { + 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) + } else { + w.WriteHeader(http.StatusNotFound) + } + case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/blobs/uploads/"): + w.Header().Set("Location", strings.TrimSuffix(r.URL.Path, "/")+"/one") + w.WriteHeader(http.StatusAccepted) + case r.Method == http.MethodPut && strings.Contains(r.URL.Path, "/blobs/uploads/"): + body, _ := readAll(r) + digest := r.URL.Query().Get("digest") + if digestOf(body) != digest { + w.WriteHeader(http.StatusBadRequest) + return + } + m.blobs[digest] = body + m.uploads++ + w.WriteHeader(http.StatusCreated) + case r.Method == http.MethodPut && strings.Contains(r.URL.Path, "/manifests/"): + body, _ := readAll(r) + m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]] = body + w.WriteHeader(http.StatusCreated) + default: + w.WriteHeader(http.StatusNotFound) + } + }) +} + +func readAll(r *http.Request) ([]byte, error) { + var buf strings.Builder + b := make([]byte, 4096) + for { + n, err := r.Body.Read(b) + buf.Write(b[:n]) + if err != nil { + break + } + } + return []byte(buf.String()), nil +} + +func TestAnUpstreamIndexIsCopiedWholeIntoTheMeshsRegistry(t *testing.T) { + src, indexDigest, srcBlobs := anUpstreamRegistry(t) + defer src.Close() + dst := &theMeshsRegistry{blobs: map[string][]byte{}, manifests: map[string][]byte{}} + dstServer := httptest.NewServer(dst.handler()) + defer dstServer.Close() + + address := strings.TrimPrefix(dstServer.URL, "http://") + r := Registry{Address: address, HTTP: src.Client()} + from := strings.TrimPrefix(src.URL, "http://") + "/library/thing:latest" + reference, err := r.MirrorImage(context.Background(), from, "hello-web/server") + if err != nil { + t.Fatal(err) + } + // Pinned by the INDEX's digest under the module's own repository: what a machine fetches is + // the whole image, whatever its architecture. + if reference != address+"/hello-web/server@"+indexDigest { + t.Fatalf("pinned as %q, not the index under the module's repository", reference) + } + // Every blob of both platforms, moved by digest, and each only once. + if len(dst.blobs) != len(srcBlobs) || dst.uploads != len(srcBlobs) { + t.Fatalf("%d of %d blobs arrived in %d uploads", len(dst.blobs), len(srcBlobs), dst.uploads) + } + for digest, body := range srcBlobs { + if string(dst.blobs[digest]) != string(body) { + t.Fatalf("blob %s did not arrive intact", digest) + } + } + // Two manifests and the index, each under its digest. + if len(dst.manifests) != 3 { + t.Fatalf("expected two manifests and an index, got %d: %v", len(dst.manifests), dst.manifests) + } + if _, ok := dst.manifests[indexDigest]; !ok { + t.Fatal("the index was not put under its digest") + } + + // Copied again, nothing is uploaded twice: blobs are content-named and already there. + if _, err := r.MirrorImage(context.Background(), from, "hello-web/server"); err != nil { + t.Fatal(err) + } + if dst.uploads != len(srcBlobs) { + t.Fatalf("a second copy uploaded blobs the registry already held: %d uploads", dst.uploads) + } +} + +func TestAReferenceIsReadTheWayARuntimeReadsIt(t *testing.T) { + for ref, want := range map[string]upstream{ + "alpine": {base: "https://registry-1.docker.io", repository: "library/alpine", reference: "latest"}, + "alpine@sha256:abc": {base: "https://registry-1.docker.io", repository: "library/alpine", reference: "sha256:abc"}, + "minio/minio:RELEASE.2025": {base: "https://registry-1.docker.io", repository: "minio/minio", reference: "RELEASE.2025"}, + "quay.io/minio/mc@sha256:def": {base: "https://quay.io", repository: "minio/mc", reference: "sha256:def"}, + "lscr.io/linuxserver/sonarr:4": {base: "https://lscr.io", repository: "linuxserver/sonarr", reference: "4"}, + "localhost:5000/x/y:1": {base: "http://localhost:5000", repository: "x/y", reference: "1"}, + } { + got, err := parseReference(ref) + if err != nil { + t.Fatalf("%s: %v", ref, err) + } + if got != want { + t.Errorf("%s: got %+v want %+v", ref, got, want) + } + } +}