The store had no working tools: its TypeScript client was never built, and nothing could count what the store holds that no record names. A Go bundle lists the store's files through its own container, reads each manifest through its door, and sets that beside the controller's records. Collection is asked of the controller, which decides and records; the store's tools never delete. The image.pushed event was declared and never emitted, and nothing consumes it, so it goes.
314 lines
9.6 KiB
Go
314 lines
9.6 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// What the store holds, read from its own files and its own door, and nothing written to either.
|
|
//
|
|
// **Why its files.** The store's door lists repositories and the tags in each, but not a manifest no tag
|
|
// names — and since the mesh pins every machine by digest, that is almost every manifest the store
|
|
// holds. The files list them all, with every blob's size. They are read through the store's own
|
|
// container, which is the only process that has them, and only read (novox/hq ADR 0251 §1).
|
|
|
|
// storageRoot is where the registry keeps its files, inside its container.
|
|
const storageRoot = "/var/lib/registry/docker/registry/v2"
|
|
|
|
// listScript lists, in three sections, every blob with its size, every manifest each repository holds,
|
|
// and what each tag names now. Busybox's find and stat (the store's image is Alpine): no -printf. An
|
|
// upload in progress is not listed, and nor are a tag's former targets.
|
|
const listScript = `set -e
|
|
cd ` + storageRoot + `
|
|
echo '#blobs'
|
|
if [ -d blobs ]; then find blobs -type f -name data | xargs -r stat -c '%s %n'; fi
|
|
echo '#revisions'
|
|
if [ -d repositories ]; then find repositories -type f -path '*/_manifests/revisions/sha256/*/link'; fi
|
|
echo '#tags'
|
|
if [ -d repositories ]; then find repositories -type f -path '*/_manifests/tags/*/current/link' -exec grep -H . {} + || true; fi
|
|
echo '#end'`
|
|
|
|
// Repo is one repository as the store's files list it.
|
|
type Repo struct {
|
|
Name string
|
|
// Revisions are the digests of every manifest it holds, tagged or not.
|
|
Revisions map[string]bool
|
|
// Tags are what each tag names now.
|
|
Tags map[string]string
|
|
}
|
|
|
|
// Listing is the store's files: every blob with its size, and every repository.
|
|
type Listing struct {
|
|
Blobs map[string]int64
|
|
Repos map[string]*Repo
|
|
}
|
|
|
|
func (l *Listing) repo(name string) *Repo {
|
|
r := l.Repos[name]
|
|
if r == nil {
|
|
r = &Repo{Name: name, Revisions: map[string]bool{}, Tags: map[string]string{}}
|
|
l.Repos[name] = r
|
|
}
|
|
return r
|
|
}
|
|
|
|
// RepoNames are the repositories in order.
|
|
func (l *Listing) RepoNames() []string {
|
|
out := make([]string, 0, len(l.Repos))
|
|
for n := range l.Repos {
|
|
out = append(out, n)
|
|
}
|
|
sort.Strings(out)
|
|
return out
|
|
}
|
|
|
|
// ParseListing reads listScript's output. A section that never ended means the listing was cut short,
|
|
// which is an error: a partial listing would make held things look absent.
|
|
func ParseListing(out string) (*Listing, error) {
|
|
l := &Listing{Blobs: map[string]int64{}, Repos: map[string]*Repo{}}
|
|
section := ""
|
|
ended := false
|
|
for _, line := range strings.Split(out, "\n") {
|
|
line = strings.TrimSpace(line)
|
|
if line == "" {
|
|
continue
|
|
}
|
|
if strings.HasPrefix(line, "#") {
|
|
section = line
|
|
ended = line == "#end"
|
|
continue
|
|
}
|
|
switch section {
|
|
case "#blobs":
|
|
size, path, ok := strings.Cut(line, " ")
|
|
n, err := strconv.ParseInt(size, 10, 64)
|
|
if !ok || err != nil {
|
|
return nil, fmt.Errorf("the store's listing has a blob line it cannot read: %q", line)
|
|
}
|
|
// blobs/sha256/ab/<hex>/data
|
|
parts := strings.Split(strings.TrimPrefix(path, "./"), "/")
|
|
if len(parts) != 5 || parts[0] != "blobs" || parts[4] != "data" {
|
|
return nil, fmt.Errorf("the store's listing has a blob at an unexpected place: %q", path)
|
|
}
|
|
l.Blobs[parts[1]+":"+parts[3]] = n
|
|
case "#revisions":
|
|
// repositories/<repo>/_manifests/revisions/sha256/<hex>/link
|
|
name, rest, ok := strings.Cut(strings.TrimPrefix(strings.TrimPrefix(line, "./"), "repositories/"), "/_manifests/revisions/")
|
|
parts := strings.Split(rest, "/")
|
|
if !ok || len(parts) != 3 || parts[2] != "link" {
|
|
return nil, fmt.Errorf("the store's listing has a manifest at an unexpected place: %q", line)
|
|
}
|
|
l.repo(name).Revisions[parts[0]+":"+parts[1]] = true
|
|
case "#tags":
|
|
// repositories/<repo>/_manifests/tags/<tag>/current/link:sha256:<hex>
|
|
path, digest, ok := strings.Cut(line, ":sha256:")
|
|
name, rest, ok2 := strings.Cut(strings.TrimPrefix(strings.TrimPrefix(path, "./"), "repositories/"), "/_manifests/tags/")
|
|
tag, _, ok3 := strings.Cut(rest, "/current/link")
|
|
if !ok || !ok2 || !ok3 || tag == "" {
|
|
return nil, fmt.Errorf("the store's listing has a tag line it cannot read: %q", line)
|
|
}
|
|
l.repo(name).Tags[tag] = "sha256:" + digest
|
|
default:
|
|
return nil, fmt.Errorf("the store's listing has a line outside any section: %q", line)
|
|
}
|
|
}
|
|
if !ended {
|
|
return nil, fmt.Errorf("the store's listing was cut short (it never reached its end)")
|
|
}
|
|
return l, nil
|
|
}
|
|
|
|
// Manifest is what one manifest names.
|
|
type Manifest struct {
|
|
MediaType string
|
|
ConfigMedia string
|
|
// Refs are the blobs it names: its configuration and its layers.
|
|
Refs []string
|
|
// Layers are its layers alone, in order.
|
|
Layers []string
|
|
// Children are the manifests an index names.
|
|
Children []string
|
|
}
|
|
|
|
// Holder is the shape of a manifest that keeps one blob: an empty configuration and one layer. The
|
|
// controller holds every archive the mesh keeps this way (novox/hq issue 253), and keeps a named
|
|
// document the same way, under a tag.
|
|
func (m *Manifest) Holder() bool {
|
|
return m != nil && m.ConfigMedia == "application/vnd.oci.empty.v1+json" && len(m.Layers) == 1
|
|
}
|
|
|
|
// ParseManifest reads a manifest's content, of any of the shapes a registry holds.
|
|
func ParseManifest(raw []byte) (*Manifest, error) {
|
|
type desc struct {
|
|
MediaType string `json:"mediaType"`
|
|
Digest string `json:"digest"`
|
|
}
|
|
var v struct {
|
|
MediaType string `json:"mediaType"`
|
|
Config *desc `json:"config"`
|
|
Layers []desc `json:"layers"`
|
|
Manifests []desc `json:"manifests"`
|
|
FSLayers []struct {
|
|
BlobSum string `json:"blobSum"`
|
|
} `json:"fsLayers"`
|
|
}
|
|
if err := json.Unmarshal(raw, &v); err != nil {
|
|
return nil, fmt.Errorf("not a manifest: %v", err)
|
|
}
|
|
m := &Manifest{MediaType: v.MediaType}
|
|
if v.Config != nil && v.Config.Digest != "" {
|
|
m.ConfigMedia = v.Config.MediaType
|
|
m.Refs = append(m.Refs, v.Config.Digest)
|
|
}
|
|
for _, d := range v.Layers {
|
|
m.Refs = append(m.Refs, d.Digest)
|
|
m.Layers = append(m.Layers, d.Digest)
|
|
}
|
|
for _, d := range v.FSLayers {
|
|
m.Refs = append(m.Refs, d.BlobSum)
|
|
m.Layers = append(m.Layers, d.BlobSum)
|
|
}
|
|
for _, d := range v.Manifests {
|
|
m.Children = append(m.Children, d.Digest)
|
|
}
|
|
return m, nil
|
|
}
|
|
|
|
var manifestAccept = []string{
|
|
"application/vnd.oci.image.manifest.v1+json",
|
|
"application/vnd.oci.image.index.v1+json",
|
|
"application/vnd.docker.distribution.manifest.v2+json",
|
|
"application/vnd.docker.distribution.manifest.list.v2+json",
|
|
"application/vnd.docker.distribution.manifest.v1+prettyjws",
|
|
}
|
|
|
|
// Store reaches the store: its files through its container, its manifests through its door.
|
|
type Store struct {
|
|
URL string
|
|
Container string
|
|
Run Runner
|
|
UID int
|
|
HTTP *http.Client
|
|
// ReadBudget bounds reading manifests in one call; what is not read in it is said.
|
|
ReadBudget time.Duration
|
|
|
|
mu sync.Mutex
|
|
// cache holds every manifest read: a manifest is named by its content's digest, so it never changes.
|
|
cache map[string]*Manifest
|
|
}
|
|
|
|
// List reads the store's files.
|
|
func (s *Store) List(ctx context.Context) (*Listing, error) {
|
|
if s.Container == "" {
|
|
return nil, fmt.Errorf("this bundle was not told the store's container (MESH_STORE_CONTAINER)")
|
|
}
|
|
out, err := docker(ctx, s.Run, s.UID, "exec", s.Container, "sh", "-c", listScript)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return ParseListing(out)
|
|
}
|
|
|
|
// Key names one manifest in one repository.
|
|
func Key(repo, digest string) string { return repo + "@" + digest }
|
|
|
|
// Manifests reads every manifest the listing names, from the cache or the door, eight at a time and
|
|
// within the read budget. Answers each one read, and each one not read with why.
|
|
func (s *Store) Manifests(ctx context.Context, l *Listing) (map[string]*Manifest, map[string]string) {
|
|
type job struct{ repo, digest string }
|
|
var jobs []job
|
|
got := map[string]*Manifest{}
|
|
unread := map[string]string{}
|
|
s.mu.Lock()
|
|
if s.cache == nil {
|
|
s.cache = map[string]*Manifest{}
|
|
}
|
|
for _, name := range l.RepoNames() {
|
|
for d := range l.Repos[name].Revisions {
|
|
if m, ok := s.cache[d]; ok {
|
|
got[Key(name, d)] = m
|
|
} else {
|
|
jobs = append(jobs, job{name, d})
|
|
}
|
|
}
|
|
}
|
|
s.mu.Unlock()
|
|
|
|
budget := s.ReadBudget
|
|
if budget == 0 {
|
|
budget = 12 * time.Second
|
|
}
|
|
ctx, cancel := context.WithTimeout(ctx, budget)
|
|
defer cancel()
|
|
var mu sync.Mutex
|
|
work := make(chan job)
|
|
var wg sync.WaitGroup
|
|
for w := 0; w < 8; w++ {
|
|
wg.Add(1)
|
|
go func() {
|
|
defer wg.Done()
|
|
for j := range work {
|
|
m, err := s.read(ctx, j.repo, j.digest)
|
|
mu.Lock()
|
|
if err != nil {
|
|
unread[Key(j.repo, j.digest)] = err.Error()
|
|
} else {
|
|
got[Key(j.repo, j.digest)] = m
|
|
}
|
|
mu.Unlock()
|
|
}
|
|
}()
|
|
}
|
|
for _, j := range jobs {
|
|
work <- j
|
|
}
|
|
close(work)
|
|
wg.Wait()
|
|
return got, unread
|
|
}
|
|
|
|
func (s *Store) read(ctx context.Context, repo, digest string) (*Manifest, error) {
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, fmt.Errorf("not read within this call's budget")
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, strings.TrimRight(s.URL, "/")+"/v2/"+repo+"/manifests/"+digest, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, a := range manifestAccept {
|
|
req.Header.Add("Accept", a)
|
|
}
|
|
client := s.HTTP
|
|
if client == nil {
|
|
client = &http.Client{Timeout: 10 * time.Second}
|
|
}
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the store's door did not answer: %v", err)
|
|
}
|
|
defer resp.Body.Close()
|
|
body, err := io.ReadAll(io.LimitReader(resp.Body, 4<<20))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if resp.StatusCode != http.StatusOK {
|
|
return nil, fmt.Errorf("the store's door answered %s", resp.Status)
|
|
}
|
|
m, err := ParseManifest(body)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
s.mu.Lock()
|
|
s.cache[digest] = m
|
|
s.mu.Unlock()
|
|
return m, nil
|
|
}
|