Files
mesh-controller/internal/builder/mirror.go
T
jschoubben 9f3790dcda Review: one local name is still a local name; a local name is unique; recovery knows it; recipes read as instructions; a tag before a digest; ask fails at once when nothing serves
A secrets object with one local name delivered no file. Two requirements could share
a local name. secret recover and the export could not tell two locals apart. The
recipe check missed continued lines and read heredoc bodies as bases. repo:tag@digest
kept the tag in the repository. ask now publishes mandatory, so a tool nothing serves
is said at once rather than after the wait.
2026-09-21 21:03:22 +02:00

324 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
}
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[:]) }