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