517 lines
17 KiB
Go
517 lines
17 KiB
Go
package offsite
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/rand"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"slices"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
// ---- fake registry ---------------------------------------------------------
|
|
|
|
type fakeRepo struct {
|
|
digests []string
|
|
tags map[string]string
|
|
}
|
|
|
|
type fakeManifest struct {
|
|
body []byte
|
|
mediaType string
|
|
}
|
|
|
|
// fakeRegistry is the registry as the sync reads it and as a restore writes
|
|
// it: one blob and manifest store shared by every repository, as
|
|
// distribution keeps it, and per-repository links.
|
|
type fakeRegistry struct {
|
|
repos map[string]*fakeRepo
|
|
blobs map[string][]byte
|
|
manifests map[string]fakeManifest
|
|
// linked records which blobs each repository holds, for a restore.
|
|
linked map[string]map[string]bool
|
|
revErr map[string]error
|
|
gone map[string]bool
|
|
corrupt map[string]bool
|
|
blobGets int
|
|
// log is what a restore did, in order.
|
|
log []string
|
|
}
|
|
|
|
func newFakeRegistry() *fakeRegistry {
|
|
return &fakeRegistry{
|
|
repos: map[string]*fakeRepo{}, blobs: map[string][]byte{}, manifests: map[string]fakeManifest{},
|
|
linked: map[string]map[string]bool{}, revErr: map[string]error{}, gone: map[string]bool{}, corrupt: map[string]bool{},
|
|
}
|
|
}
|
|
|
|
func digestOf(b []byte) string {
|
|
sum := sha256.Sum256(b)
|
|
return "sha256:" + hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
func (r *fakeRegistry) blob(data []byte) map[string]any {
|
|
d := digestOf(data)
|
|
r.blobs[d] = data
|
|
return map[string]any{"mediaType": "application/octet-stream", "digest": d, "size": len(data)}
|
|
}
|
|
|
|
func (r *fakeRegistry) link(repo, digest string, tags ...string) {
|
|
rp := r.repos[repo]
|
|
if rp == nil {
|
|
rp = &fakeRepo{tags: map[string]string{}}
|
|
r.repos[repo] = rp
|
|
}
|
|
if !slices.Contains(rp.digests, digest) {
|
|
rp.digests = append(rp.digests, digest)
|
|
}
|
|
for _, t := range tags {
|
|
rp.tags[t] = digest
|
|
}
|
|
}
|
|
|
|
func (r *fakeRegistry) putManifest(mediaType string, m map[string]any) string {
|
|
m["schemaVersion"] = 2
|
|
m["mediaType"] = mediaType
|
|
body, _ := json.Marshal(m)
|
|
d := digestOf(body)
|
|
r.manifests[d] = fakeManifest{body: body, mediaType: mediaType}
|
|
return d
|
|
}
|
|
|
|
// image stores an image manifest of the given layers in repo and returns its digest.
|
|
func (r *fakeRegistry) image(repo, tag string, layers ...[]byte) string {
|
|
var ls []any
|
|
for _, l := range layers {
|
|
ls = append(ls, r.blob(l))
|
|
}
|
|
cfg := r.blob([]byte(fmt.Sprintf(`{"architecture":"arm64","repo":%q,"tag":%q}`, repo, tag)))
|
|
d := r.putManifest(mediaDockerManifest2, map[string]any{"config": cfg, "layers": ls})
|
|
r.link(repo, d, tag)
|
|
return d
|
|
}
|
|
|
|
// index stores an image index over children in repo.
|
|
func (r *fakeRegistry) index(repo, tag string, children ...string) string {
|
|
var ms []any
|
|
for _, c := range children {
|
|
ms = append(ms, map[string]any{"mediaType": mediaOCIManifest, "digest": c, "size": len(r.manifests[c].body)})
|
|
}
|
|
d := r.putManifest(mediaOCIIndex, map[string]any{"manifests": ms})
|
|
r.link(repo, d, tag)
|
|
return d
|
|
}
|
|
|
|
func (r *fakeRegistry) Repositories(context.Context) ([]string, error) {
|
|
var out []string
|
|
for name := range r.repos {
|
|
out = append(out, name)
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
func (r *fakeRegistry) Revisions(_ context.Context, repo string) ([]string, map[string]string, error) {
|
|
if err := r.revErr[repo]; err != nil {
|
|
return nil, nil, err
|
|
}
|
|
rp := r.repos[repo]
|
|
return slices.Clone(rp.digests), rp.tags, nil
|
|
}
|
|
|
|
func (r *fakeRegistry) Manifest(_ context.Context, repo, digest string) ([]byte, string, error) {
|
|
m, ok := r.manifests[digest]
|
|
if !ok || r.gone[digest] {
|
|
return nil, "", fmt.Errorf("%w: manifest %s", ErrImageGone, digest)
|
|
}
|
|
return m.body, m.mediaType, nil
|
|
}
|
|
|
|
func (r *fakeRegistry) Blob(_ context.Context, repo, digest string) (io.ReadCloser, error) {
|
|
r.blobGets++
|
|
b, ok := r.blobs[digest]
|
|
if !ok || r.gone[digest] {
|
|
return nil, fmt.Errorf("%w: blob %s", ErrImageGone, digest)
|
|
}
|
|
if r.corrupt[digest] {
|
|
b = append(bytes.Clone(b), 'x')
|
|
}
|
|
return io.NopCloser(bytes.NewReader(b)), nil
|
|
}
|
|
|
|
func (r *fakeRegistry) PutBlob(_ context.Context, repo, digest string, size int64, open func() (io.ReadCloser, error)) error {
|
|
if r.linked[repo][digest] {
|
|
return nil
|
|
}
|
|
rc, err := open()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer rc.Close()
|
|
data, err := io.ReadAll(rc)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if int64(len(data)) != size || digestOf(data) != digest {
|
|
return fmt.Errorf("blob upload does not match %s", digest)
|
|
}
|
|
r.blobs[digest] = data
|
|
if r.linked[repo] == nil {
|
|
r.linked[repo] = map[string]bool{}
|
|
}
|
|
r.linked[repo][digest] = true
|
|
r.log = append(r.log, "blob "+repo+" "+digest)
|
|
return nil
|
|
}
|
|
|
|
func (r *fakeRegistry) PutManifest(_ context.Context, repo, reference, mediaType string, body []byte) error {
|
|
d := digestOf(body)
|
|
if strings.HasPrefix(reference, "sha256:") && reference != d {
|
|
return fmt.Errorf("manifest %s hashes to %s", reference, d)
|
|
}
|
|
// Like distribution: everything a manifest names must be in the repository.
|
|
im, err := describeManifest(body, mediaType)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, b := range im.Blobs {
|
|
if !r.linked[repo][b.Digest] {
|
|
return fmt.Errorf("blob unknown: %s", b.Digest)
|
|
}
|
|
}
|
|
for _, c := range im.Children {
|
|
if r.repos[repo] == nil || !slices.Contains(r.repos[repo].digests, c) {
|
|
return fmt.Errorf("manifest unknown: %s", c)
|
|
}
|
|
}
|
|
r.manifests[d] = fakeManifest{body: body, mediaType: mediaType}
|
|
if strings.HasPrefix(reference, "sha256:") {
|
|
r.link(repo, d)
|
|
} else {
|
|
r.link(repo, d, reference)
|
|
}
|
|
r.log = append(r.log, "manifest "+repo+" "+reference)
|
|
return nil
|
|
}
|
|
|
|
func layer(n int) []byte {
|
|
b := make([]byte, n)
|
|
rand.Read(b)
|
|
return b
|
|
}
|
|
|
|
func newImageSyncer(t *testing.T, reg *fakeRegistry) (*Syncer, *memBucket, *time.Time) {
|
|
t.Helper()
|
|
b := newMemBucket()
|
|
clock := now
|
|
return &Syncer{
|
|
Bucket: b, Catalog: &fakeCatalog{}, Key: testKey(t), Images: reg,
|
|
Now: func() time.Time { return clock },
|
|
}, b, &clock
|
|
}
|
|
|
|
func keysUnder(b *memBucket, prefix string) []string {
|
|
objs, _ := b.List(context.Background(), prefix)
|
|
var out []string
|
|
for _, o := range objs {
|
|
out = append(out, o.Key)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// ---- tests -----------------------------------------------------------------
|
|
|
|
// TestSyncImagesCopiesUserImages: user repositories are copied blob by blob,
|
|
// each blob once however many images share it, and the platform's reserved
|
|
// repositories are left alone. A second run with nothing new sends nothing.
|
|
func TestSyncImagesCopiesUserImages(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
shared := layer(3000)
|
|
a := reg.image("user-uploads/sub-a", "latest", shared, layer(70000))
|
|
child := reg.image("e2e/multi", "arm64", shared)
|
|
idx := reg.index("e2e/multi", "latest", child)
|
|
reg.image("felis/felis", "v1", layer(5000))
|
|
reg.image("mirror/trivy-db", "2", layer(5000))
|
|
|
|
s, b, _ := newImageSyncer(t, reg)
|
|
res, err := s.Run(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("Run: %v", err)
|
|
}
|
|
if res.ImageRepos != 2 || res.Images != 3 || res.ImagesUploaded != 3 {
|
|
t.Fatalf("result = %+v, want 2 repositories, 3 images, 3 uploaded", res)
|
|
}
|
|
// shared layer once, one unique layer, two configs
|
|
if got := keysUnder(b, imageBlobsDir); len(got) != 4 || res.ImageBlobsUploaded != 4 {
|
|
t.Fatalf("blobs in bucket = %v (uploaded %d), want 4", got, res.ImageBlobsUploaded)
|
|
}
|
|
for _, d := range []string{a, child, idx} {
|
|
if _, ok := b.objs[imageManifestKey(d)]; !ok {
|
|
t.Errorf("manifest %s not in the bucket", d)
|
|
}
|
|
}
|
|
x, err := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, ok := x.Repositories["felis/felis"]; ok {
|
|
t.Fatal("reserved repository felis/felis was copied")
|
|
}
|
|
if got := x.Repositories["e2e/multi"]; got.Tags["latest"] != idx || got.Tags["arm64"] != child {
|
|
t.Fatalf("e2e/multi tags = %v", got.Tags)
|
|
}
|
|
|
|
puts, gets := b.puts, reg.blobGets
|
|
res, err = s.Run(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("second Run: %v", err)
|
|
}
|
|
if b.puts != puts || reg.blobGets != gets {
|
|
t.Fatalf("second run sent %d objects and read %d blobs, want none", b.puts-puts, reg.blobGets-gets)
|
|
}
|
|
if len(keysUnder(b, imageIndexDir)) != 1 {
|
|
t.Fatalf("an unchanged registry wrote another index: %v", keysUnder(b, imageIndexDir))
|
|
}
|
|
}
|
|
|
|
// TestSyncImagesCopiesPinnedPlatformRevisions: of felis/ and mirror/ only the
|
|
// revisions a pin names are copied, without tags; a restore puts them back by
|
|
// digest and leaves the tags the installer pushed on the new host alone. When
|
|
// the pins cannot be read, the previous copy of those revisions is kept and
|
|
// the run fails.
|
|
func TestSyncImagesCopiesPinnedPlatformRevisions(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
old := reg.image("felis/paper", "demo", layer(4000))
|
|
current := reg.image("felis/paper", "demo", layer(4000))
|
|
reg.image("mirror/trivy", "1", layer(2000))
|
|
user := reg.image("user/a", "v1", layer(1000))
|
|
|
|
s, b, clock := newImageSyncer(t, reg)
|
|
s.ImagePins = func(context.Context) (map[string][]string, error) {
|
|
return map[string][]string{"felis/paper": {old}, "user/a": {user}}, nil
|
|
}
|
|
res, err := s.Run(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("Run: %v", err)
|
|
}
|
|
x, err := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
paper, ok := x.Repositories["felis/paper"]
|
|
if !ok || !slices.Equal(paper.Manifests, []string{old}) || len(paper.Tags) != 0 {
|
|
t.Fatalf("felis/paper = %+v, want only the pinned %s and no tags", paper, old)
|
|
}
|
|
if _, ok := x.Repositories["mirror/trivy"]; ok {
|
|
t.Fatal("unpinned mirror/trivy was copied")
|
|
}
|
|
if got := x.Repositories["user/a"]; got.Tags["v1"] != user {
|
|
t.Fatalf("user/a = %+v", got)
|
|
}
|
|
|
|
fresh := newFakeRegistry()
|
|
rebuilt := fresh.image("felis/paper", "demo", layer(4000))
|
|
if _, err := FetchImages(context.Background(), b, s.Key, x, fresh, nil); err != nil {
|
|
t.Fatalf("FetchImages: %v", err)
|
|
}
|
|
if got := fresh.repos["felis/paper"]; got.tags["demo"] != rebuilt || !slices.Contains(got.digests, old) {
|
|
t.Fatalf("restored felis/paper: tags %v digests %v, want demo=%s and %s present", got.tags, got.digests, rebuilt, old)
|
|
}
|
|
if slices.Contains(fresh.repos["felis/paper"].digests, current) {
|
|
t.Fatal("the unpinned revision was restored")
|
|
}
|
|
|
|
*clock = clock.Add(time.Hour)
|
|
s.ImagePins = func(context.Context) (map[string][]string, error) { return nil, errors.New("cluster down") }
|
|
res, err = s.Run(context.Background())
|
|
if err == nil {
|
|
t.Fatal("Run succeeded without the pins")
|
|
}
|
|
x, err = LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if got := x.Repositories["felis/paper"]; !slices.Equal(got.Manifests, []string{old}) {
|
|
t.Fatalf("without the pins felis/paper = %+v, want the previous copy kept", got)
|
|
}
|
|
}
|
|
|
|
// TestSyncImagesRejectsCorruptBlob: bytes that do not hash to the digest never
|
|
// become that digest's object, the image stays out of the index, and the run
|
|
// fails so the next one tries again.
|
|
func TestSyncImagesRejectsCorruptBlob(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
bad := layer(100)
|
|
reg.image("user/a", "v1", bad)
|
|
reg.corrupt[digestOf(bad)] = true
|
|
|
|
s, b, _ := newImageSyncer(t, reg)
|
|
res, err := s.Run(context.Background())
|
|
if err == nil || !strings.Contains(err.Error(), "longer than") {
|
|
t.Fatalf("Run = %v, want the corrupt blob to fail it", err)
|
|
}
|
|
if _, ok := b.objs[imageBlobKey(digestOf(bad))]; ok {
|
|
t.Fatal("the corrupt blob was stored")
|
|
}
|
|
if len(keysUnder(b, imageManifestsDir)) != 0 || res.Images != 0 {
|
|
t.Fatalf("manifest stored without its blob: %v, images=%d", keysUnder(b, imageManifestsDir), res.Images)
|
|
}
|
|
}
|
|
|
|
// TestSyncImagesSkipsWhatTheRegistryLost: a manifest whose blob the registry no
|
|
// longer has is reported, not copied, and does not fail the run; one copied
|
|
// earlier keeps its off-site copy.
|
|
func TestSyncImagesSkipsWhatTheRegistryLost(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
keptLayer := layer(100)
|
|
kept := reg.image("user/a", "v1", keptLayer)
|
|
s, b, _ := newImageSyncer(t, reg)
|
|
if _, err := s.Run(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
lostLayer := layer(100)
|
|
lost := reg.image("user/a", "v2", lostLayer)
|
|
reg.gone[digestOf(lostLayer)] = true
|
|
res, err := s.Run(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("Run = %v, want a lost blob to be reported only", err)
|
|
}
|
|
if len(res.ImagesIncomplete) != 1 || !strings.Contains(res.ImagesIncomplete[0], lost) {
|
|
t.Fatalf("incomplete = %v, want %s", res.ImagesIncomplete, lost)
|
|
}
|
|
x, _ := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex)
|
|
if got := x.Repositories["user/a"]; !slices.Equal(got.Manifests, []string{kept}) || got.Tags["v2"] != "" {
|
|
t.Fatalf("user/a = %+v, want only the complete image", got)
|
|
}
|
|
}
|
|
|
|
// TestSyncImagesKeepsReplacedStates: a registry that comes back empty (a host
|
|
// being rebuilt) writes a new index version but prunes nothing; the objects of
|
|
// the old state go only ImageHistory after it was replaced.
|
|
func TestSyncImagesKeepsReplacedStates(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
reg.image("user/a", "v1", layer(100))
|
|
s, b, clock := newImageSyncer(t, reg)
|
|
if _, err := s.Run(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
objects := len(keysUnder(b, imageBlobsDir)) + len(keysUnder(b, imageManifestsDir))
|
|
|
|
reg.repos = map[string]*fakeRepo{}
|
|
*clock = clock.Add(time.Hour)
|
|
res, err := s.Run(context.Background())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if res.ImageRepos != 0 || len(keysUnder(b, imageIndexDir)) != 2 || res.ImageObjectsPruned != 0 {
|
|
t.Fatalf("after emptying: repos=%d versions=%v pruned=%d", res.ImageRepos, keysUnder(b, imageIndexDir), res.ImageObjectsPruned)
|
|
}
|
|
if got := len(keysUnder(b, imageBlobsDir)) + len(keysUnder(b, imageManifestsDir)); got != objects {
|
|
t.Fatalf("objects = %d, want the old state's %d kept", got, objects)
|
|
}
|
|
// A restore now refuses the empty newest state and names the old one.
|
|
if _, _, err := ChooseImageIndex(context.Background(), b, s.Key, ""); err == nil || !strings.Contains(err.Error(), "-at") {
|
|
t.Fatalf("ChooseImageIndex = %v, want it to point at -at", err)
|
|
}
|
|
|
|
*clock = clock.Add(ImageHistory - 2*time.Hour)
|
|
if res, _ := s.Run(context.Background()); res.ImageObjectsPruned != 0 {
|
|
t.Fatalf("pruned %d objects inside the history window", res.ImageObjectsPruned)
|
|
}
|
|
*clock = clock.Add(3 * time.Hour)
|
|
res, err = s.Run(context.Background())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if res.ImageObjectsPruned != objects+1 || len(keysUnder(b, imageBlobsDir)) != 0 || len(keysUnder(b, imageIndexDir)) != 1 {
|
|
t.Fatalf("pruned %d, left blobs %v and versions %v", res.ImageObjectsPruned, keysUnder(b, imageBlobsDir), keysUnder(b, imageIndexDir))
|
|
}
|
|
}
|
|
|
|
// TestSyncImagesUnreadableRepoKeepsItsEntry: a repository the registry fails to
|
|
// list keeps its last recorded state, and the failed run prunes nothing.
|
|
func TestSyncImagesUnreadableRepoKeepsItsEntry(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
a := reg.image("user/a", "v1", layer(100))
|
|
s, b, clock := newImageSyncer(t, reg)
|
|
if _, err := s.Run(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
reg.revErr["user/a"] = errors.New("gate restarting")
|
|
reg.image("user/b", "v1", layer(100))
|
|
*clock = clock.Add(ImageHistory * 2)
|
|
res, err := s.Run(context.Background())
|
|
if err == nil {
|
|
t.Fatal("Run succeeded with an unreadable repository")
|
|
}
|
|
x, _ := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex)
|
|
if got := x.Repositories["user/a"]; !slices.Equal(got.Manifests, []string{a}) {
|
|
t.Fatalf("user/a = %+v, want its last recorded state", got)
|
|
}
|
|
if _, ok := x.Repositories["user/b"]; !ok || res.ImageObjectsPruned != 0 {
|
|
t.Fatalf("user/b missing or pruned %d", res.ImageObjectsPruned)
|
|
}
|
|
}
|
|
|
|
// TestFetchImagesRestoresRegistry: a fresh registry gets back every user image
|
|
// with the same digests and tags, blobs before the manifests that name them and
|
|
// an index after its children; a second run pushes nothing.
|
|
func TestFetchImagesRestoresRegistry(t *testing.T) {
|
|
reg := newFakeRegistry()
|
|
shared := layer(3000)
|
|
a := reg.image("user/a", "v1", shared, layer(9000))
|
|
child := reg.image("user/multi", "arm64", shared)
|
|
idx := reg.index("user/multi", "latest", child)
|
|
s, b, _ := newImageSyncer(t, reg)
|
|
if _, err := s.Run(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
stamp, x, err := ChooseImageIndex(context.Background(), b, s.Key, "")
|
|
if err != nil || stamp == "" {
|
|
t.Fatalf("ChooseImageIndex = %q, %v", stamp, err)
|
|
}
|
|
fresh := newFakeRegistry()
|
|
res, err := FetchImages(context.Background(), b, s.Key, x, fresh, nil)
|
|
if err != nil {
|
|
t.Fatalf("FetchImages: %v", err)
|
|
}
|
|
if res.Repositories != 2 || res.Manifests != 3 || res.Tags != 3 || res.BlobsPushed != 5 {
|
|
t.Fatalf("result = %+v", res)
|
|
}
|
|
for repo, want := range map[string]map[string]string{
|
|
"user/a": {"v1": a},
|
|
"user/multi": {"arm64": child, "latest": idx},
|
|
} {
|
|
for tag, d := range want {
|
|
if got := fresh.repos[repo].tags[tag]; got != d {
|
|
t.Errorf("%s:%s = %s, want %s", repo, tag, got, d)
|
|
}
|
|
}
|
|
}
|
|
if !bytes.Equal(fresh.manifests[idx].body, reg.manifests[idx].body) {
|
|
t.Fatal("restored index differs from the original")
|
|
}
|
|
pos := func(entry string) int { return slices.Index(fresh.log, entry) }
|
|
if pos("manifest user/multi "+child) > pos("manifest user/multi "+idx) {
|
|
t.Fatalf("index pushed before its child: %v", fresh.log)
|
|
}
|
|
|
|
n := len(fresh.log)
|
|
if _, err := FetchImages(context.Background(), b, s.Key, x, fresh, nil); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
for _, e := range fresh.log[n:] {
|
|
if strings.HasPrefix(e, "blob ") {
|
|
t.Fatalf("second restore uploaded a blob again: %s", e)
|
|
}
|
|
}
|
|
}
|