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) + } + } +}