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.
324 lines
12 KiB
Go
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[:]) }
|