fix(offsite): 服务器或白名单按摘要钉住的 felis/ 与 mirror/ 版本也进异地副本,恢复只按摘要推回不动安装器的 tag

This commit is contained in:
Lemon-miaow committed 2026-09-25 05:00:02 +08:00
1 parent 8fb3d298ae
commit 2779d8f5cf
8 files changed
+158 -20

No files matched your search

+32 -12
View File
@@ -1,10 +1,13 @@
package offsite
// The off-site copy of the platform registry's user images. The installer
// pushes everything under felis/ and mirror/ again on any host, so those are
// left out; every other repository holds builds that exist nowhere else: a
// lost registry volume would otherwise leave each server pinned to one of them
// in ImagePullBackOff until someone rebuilds it.
// The off-site copy of the platform registry's user images. Every repository
// outside felis/ and mirror/ holds builds that exist nowhere else: a lost
// registry volume would otherwise leave each server pinned to one of them in
// ImagePullBackOff until someone rebuilds it. The installer pushes felis/ and
// mirror/ again on any host, but under the same tags at new digests (a build
// is not reproducible, an upstream tag moves on), so from those only the
// revisions a server or a whitelist entry pins by digest are copied, without
// their tags: the tags belong to whatever the installer pushed last.
//
// The bucket holds each blob and manifest once, by digest, encrypted like
// everything else, plus an index that says which repository holds which
@@ -234,19 +237,28 @@ func (s *Syncer) syncImages(ctx context.Context, res *Result, fail func(string,
return
}
sort.Strings(repos)
listed := map[string]bool{}
pins, pinsKnown := map[string][]string{}, true
if s.ImagePins != nil {
if pins, err = s.ImagePins(ctx); err != nil {
fail("list the images servers and the whitelist pin: %v", err)
pinsKnown, c.failed = false, true
}
}
for _, repo := range repos {
if reservedRepo(repo) {
if reservedRepo(repo) && !pinsKnown {
c.carry(repo)
continue
}
if reservedRepo(repo) && len(pins[repo]) == 0 {
continue
}
listed[repo] = true
if ctx.Err() != nil {
c.fail("stopped before %s: %v", repo, ctx.Err())
c.failed = true
c.carry(repo)
continue
}
c.repo(ctx, repo)
c.repo(ctx, repo, pins[repo])
}
stamp := ""
@@ -277,9 +289,10 @@ func (s *Syncer) syncImages(ctx context.Context, res *Result, fail func(string,
}
}
// repo copies one repository into c.next. A repository that cannot be read in
// full keeps the entry the previous index had for it.
func (c *imageCopy) repo(ctx context.Context, repo string) {
// repo copies one repository into c.next; of a reserved one, only the pinned
// revisions and none of its tags. A repository that cannot be read in full
// keeps the entry the previous index had for it.
func (c *imageCopy) repo(ctx context.Context, repo string, pinned []string) {
digests, tags, err := c.s.Images.Revisions(ctx, repo)
if err != nil {
c.fail("list the images of %s: %v", repo, err)
@@ -287,6 +300,10 @@ func (c *imageCopy) repo(ctx context.Context, repo string) {
c.carry(repo)
return
}
if reservedRepo(repo) {
digests = slices.DeleteFunc(digests, func(d string) bool { return !slices.Contains(pinned, d) })
tags = nil
}
entry := ImageRepo{Tags: map[string]string{}}
short := false
for _, d := range digests {
@@ -732,6 +749,9 @@ func FetchImages(ctx context.Context, b Bucket, key []byte, x *ImageIndex, t Ima
pushedBlobs := map[string]bool{}
for _, repo := range slices.Sorted(maps.Keys(x.Repositories)) {
entry := x.Repositories[repo]
if reservedRepo(repo) {
entry.Tags = nil // the installer's tags stay where it put them
}
done := map[string]bool{}
var push func(d string) error
push = func(d string) error {
+62
View File
@@ -280,6 +280,68 @@ func TestSyncImagesCopiesUserImages(t *testing.T) {
}
}
// 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.
+4
View File
@@ -95,6 +95,10 @@ type Syncer struct {
// Images is the platform registry whose user images are copied (images.go);
// nil copies none.
Images ImageSource
// ImagePins lists, per repository, the digests a server's spec or a
// whitelist entry pins. Under felis/ and mirror/ only those revisions are
// copied (images.go); nil copies none there.
ImagePins func(ctx context.Context) (map[string][]string, error)
// UploadsDir is the host directory of the uploads volume, whose submission
// contexts are copied (uploads.go); empty copies none.
UploadsDir string
+3 -3
View File
@@ -172,7 +172,7 @@ func (p *Pruner) log() *slog.Logger {
func Plan(indexes map[string]*registrygate.Index, refs []string, host string, now time.Time, grace time.Duration, keepTagged int) []Target {
keep := map[Target]bool{}
for _, ref := range refs {
repo, tag, digest, ok := parseRef(ref, host)
repo, tag, digest, ok := ParseRef(ref, host)
if !ok {
continue
}
@@ -255,9 +255,9 @@ func reserved(repo string) bool {
return false
}
// parseRef splits host/repo[:tag][@digest] for refs under host. A ref with
// ParseRef splits host/repo[:tag][@digest] for refs under host. A ref with
// neither tag nor digest means :latest, as it does for every image client.
func parseRef(ref, host string) (repo, tag, digest string, ok bool) {
func ParseRef(ref, host string) (repo, tag, digest string, ok bool) {
rest, ok := strings.CutPrefix(strings.TrimSpace(ref), host+"/")
if !ok || rest == "" {
return "", "", "", false
+2 -2
View File
@@ -124,9 +124,9 @@ func TestParseRef(t *testing.T) {
{"other:5000/felis/paper:demo", "", "", "", false},
{host + "/", "", "", "", false},
} {
repo, tag, digest, ok := parseRef(c.ref, host)
repo, tag, digest, ok := ParseRef(c.ref, host)
if repo != c.repo || tag != c.tag || digest != c.digest || ok != c.ok {
t.Errorf("parseRef(%q) = %q %q %q %v, want %q %q %q %v", c.ref, repo, tag, digest, ok, c.repo, c.tag, c.digest, c.ok)
t.Errorf("ParseRef(%q) = %q %q %q %v, want %q %q %q %v", c.ref, repo, tag, digest, ok, c.repo, c.tag, c.digest, c.ok)
}
}
}