Files
mesh-controller/internal/builder/mirror.go
T
jschoubben 338d033632 A manifest HEAD says what it accepts, or the registry answers 404
The check that skips copying a base the mesh already holds asked with no Accept header, and a
registry answers a manifest only in a media type the caller named: the same digest answered 200 with
the manifest types and 404 without them. So the builder concluded it held nothing, copied every
vendor base again, and exhausted the public hub's pull limit a second time today.

The test could not have caught it, because the fake registry answered a manifest HEAD regardless of
Accept — more permissive than the thing it stands in for. It is now as strict as a real registry, and
fails without the fix.
2026-09-28 13:02:08 +02:00

338 lines
12 KiB
Go

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:]
// `repo:tag@digest` is what a runtime prints; the digest names the image and the tag is
// only what it was called. The tag is not part of the repository.
if colon := strings.LastIndex(name, ":"); colon > strings.LastIndex(name, "/") {
name = name[:colon]
}
} 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
}
// **Already held is already mirrored.** A base is named by digest, and a digest this registry
// holds under the module's repository is the same bytes whatever upstream would say — so
// upstream is not asked. Asked every build, the public hub's anonymous pull limit was reached
// 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)
if err != nil {
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
}
if held {
return r.Address + "/" + repository + "@" + where.reference, nil
}
}
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[:]) }