Files
mesh-controller/internal/builder/mirror.go
T
jochen 80a928e4c1
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
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).
2026-10-08 14:51:13 +02:00

445 lines
17 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)
}
// BaseMirrorer copies a base a build stands on into the one repository the mesh keeps for that
// upstream image, taking what an earlier copy under `former` already holds rather than asking
// upstream for it again (novox/hq ADR 0257).
type BaseMirrorer interface {
MirrorBase(ctx context.Context, from, repository, former string) (string, error)
}
// MirrorPrefix is the namespace of the repositories a base is mirrored into: `upstream/<host>/<path>`,
// one repository per upstream image, whichever modules stand on it (novox/hq ADR 0257).
const MirrorPrefix = "upstream/"
// MirrorRepository is the repository the mesh keeps an upstream image's copy in: the image's own
// host and path under MirrorPrefix, so two modules standing on one image stand on one copy, and two
// images that share a path on different hosts are never confused. Lower case and without a port's
// colon, because a repository name allows neither.
func MirrorRepository(from string) (string, error) {
where, err := parseReference(from)
if err != nil {
return "", err
}
host := strings.TrimPrefix(strings.TrimPrefix(where.base, "https://"), "http://")
if host == "registry-1.docker.io" {
host = "docker.io"
}
host = strings.ReplaceAll(host, ":", "-")
return strings.ToLower(MirrorPrefix + host + "/" + where.repository), nil
}
// FormerMirrorRepository is where a module's base was copied before ADR 0257: under the module's own
// repository, one copy per module (ADR 0097). Read only, as a source of what is already held.
func FormerMirrorRepository(module, arg string) string {
return module + "/on-" + strings.ToLower(arg)
}
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
// mountFrom is a repository of the mesh's own registry the source is, so a blob is mounted from
// it rather than read and written again.
mountFrom 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) {
return r.MirrorBase(ctx, from, repository, "")
}
// MirrorBase is MirrorImage, and where `former` — a repository of this registry — already holds the
// image by its digest, the copy is made from there: manifests read from this registry, blobs mounted
// across rather than moved (novox/hq ADR 0257). Upstream is asked only for what no repository here
// holds, which is what kept the public hub's anonymous pull limit out of reach when every module's
// base moved into one shared repository.
func (r Registry) MirrorBase(ctx context.Context, from, repository, former 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 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)
}
if held {
return r.Address + "/" + repository + "@" + where.reference, nil
}
}
src := &source{client: r.client()}
if former != "" && former != repository && strings.HasPrefix(where.reference, "sha256:") {
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)
}
if held {
// The same bytes, already here: copied from this registry, each blob mounted.
where = upstream{base: "http://" + r.Address, repository: former, reference: where.reference}
src.mountFrom = former
}
}
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
}
// 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.
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
}
// **Mounted where this registry already holds it** (novox/hq ADR 0257): a blob is stored once
// whichever repositories link it, so a mount moves no bytes. A registry that cannot mount answers
// with an ordinary upload's location, and the blob is moved as before.
uploads := base + "/blobs/uploads/"
if src.mountFrom != "" {
uploads += "?mount=" + digest + "&from=" + src.mountFrom
}
start, err := http.NewRequestWithContext(ctx, http.MethodPost, 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.StatusCreated && src.mountFrom != "" {
return nil
}
if begun.StatusCode != http.StatusAccepted {
return fmt.Errorf("%s answered %s when asked where to put a blob", base, begun.Status)
}
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)
}
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[:]) }