Loading internal/registrygate/index.go +43 −2 Changes for internal/registrygate/index.go: 43 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -3,6 +3,7 @@ package registrygate import ( "encoding/json" "errors" "io" "io/fs" "net/http" "os" Loading @@ -15,7 +16,7 @@ import ( // IndexPathPrefix serves the manifest index of one repository: // // GET /felis/manifests/<repo> {"revisions":[{"digest":…,"pushed":…}],"tags":{"<tag>":"<digest>"}} // GET /felis/manifests/<repo> {"revisions":[{"digest":…,"pushed":…,"children":[…]}],"tags":{"<tag>":"<digest>"}} // // The registry API can list a repository's tags but not its manifests, so a // manifest a tag moved off (every rebuild of :demo or :latest leaves one) is Loading @@ -39,6 +40,45 @@ type Index struct { type Revision struct { Digest string `json:"digest"` Pushed time.Time `json:"pushed"` // Children lists the manifests an image index (or Docker manifest list) names: // one per platform, plus BuildKit's attestation manifests. Each is a revision // of its own, untagged and never spelled in an image ref, yet a pull of the // index fetches them, so whoever keeps the index must keep them too. Empty for // a single-platform manifest. Children []string `json:"children,omitempty"` } // maxManifestBytes caps how much of a revision's blob is read to find its // children. It is the registry's own limit on a pushed manifest. const maxManifestBytes = 4 << 20 // manifestChildren reads a manifest blob off the filesystem driver's layout and // returns the digests it names when it is an index. A blob that is missing, // oversized or unparsable names nothing: a pull of it fails already, and the // pruner then treats it like any other revision. func manifestChildren(root, hex string) []string { f, err := os.Open(filepath.Join(root, "docker", "registry", "v2", "blobs", "sha256", hex[:2], hex, "data")) if err != nil { return nil } defer f.Close() // The manifests array is what makes an index, whatever mediaType says (an OCI // index may omit it); a single-platform manifest has layers and no manifests. var m struct { Manifests []struct { Digest string `json:"digest"` } `json:"manifests"` } if err := json.NewDecoder(io.LimitReader(f, maxManifestBytes)).Decode(&m); err != nil { return nil } var out []string for _, c := range m.Manifests { if digestRE.MatchString(c.Digest) { out = append(out, c.Digest) } } return out } // repoNameRE is the distribution reference grammar for a repository path. Every Loading Loading @@ -96,7 +136,8 @@ func ReadIndex(root, repo string) (*Index, error) { if err != nil { return nil, err } idx.Revisions = append(idx.Revisions, Revision{Digest: "sha256:" + e.Name(), Pushed: st.ModTime().UTC()}) idx.Revisions = append(idx.Revisions, Revision{Digest: "sha256:" + e.Name(), Pushed: st.ModTime().UTC(), Children: manifestChildren(root, e.Name())}) } sort.Slice(idx.Revisions, func(i, j int) bool { return idx.Revisions[i].Digest < idx.Revisions[j].Digest }) Loading internal/registryprune/prune.go +29 −4 Changes for internal/registryprune/prune.go: 29 added lines, 4 removed lines. Original line number Diff line number Diff line Loading @@ -4,7 +4,8 @@ // previous manifest behind, a whitelist entry an admin removes keeps its image, // and each installer run adds a platform image. // // A manifest is kept when any of these hold: // A manifest is kept when any of these hold, and so is every manifest a kept // image index names (its per-platform images and attestations): // // - an image reference the platform still depends on names it: a whitelist // entry (by tag, by wildcard tag, or by digest), a server's spec (servers are Loading Loading @@ -211,16 +212,40 @@ func Plan(indexes map[string]*registrygate.Index, refs []string, host string, no } } for _, r := range idx.Revisions { t := Target{repo, r.Digest} if keep[t] || now.Sub(r.Pushed) < grace { continue if now.Sub(r.Pushed) < grace { keep[Target{repo, r.Digest}] = true } } keepChildren(repo, idx, keep) for _, r := range idx.Revisions { if t := (Target{repo, r.Digest}); !keep[t] { out = append(out, t) } } } return out } // keepChildren extends keep from each kept image index to the manifests it names: // the per-platform images and their attestations, which no ref spells but every // pull of the index fetches. An index can name another index, so it runs to a // fixed point. func keepChildren(repo string, idx *registrygate.Index, keep map[Target]bool) { for grew := true; grew; { grew = false for _, r := range idx.Revisions { if !keep[Target{repo, r.Digest}] { continue } for _, c := range r.Children { if t := (Target{repo, c}); !keep[t] { keep[t], grew = true, true } } } } } // newestTagged returns the n most recently pushed distinct digests any tag in idx // names. func newestTagged(idx *registrygate.Index, pushed map[string]time.Time, n int) []string { Loading internal/registryprune/prune_test.go +104 −0 Changes for internal/registryprune/prune_test.go: 104 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -111,6 +111,40 @@ func TestPlanKeepsTheNewestTaggedPlatformImages(t *testing.T) { } } // A BuildKit push stores an image index plus one untagged revision per platform // and per attestation, and a server's spec pins the index. The children must stay // with it: the VM's test-one lost its arm64 manifest this way and could no longer // pull the paper image it was pinned to. func TestPlanKeepsTheManifestsAKeptIndexNames(t *testing.T) { withKids := func(r registrygate.Revision, kids ...string) registrygate.Revision { for _, k := range kids { r.Children = append(r.Children, dg(k)) } return r } indexes := map[string]*registrygate.Index{ "felis/paper": { Revisions: []registrygate.Revision{ // i: index pinned by a server, naming platform p and attestation q. withKids(rev("i", 30*day), "p", "q"), rev("p", 30*day), rev("q", 30*day), // j: an unused index; it goes with its children r and s. withKids(rev("j", 40*day), "r", "s"), rev("r", 40*day), rev("s", 40*day), // k: the newest tagged index names a nested index n, which names m. withKids(rev("k", 2*day), "n"), withKids(rev("n", 2*day), "m"), rev("m", 2*day), // o: a child of an index the registry no longer holds: goes. rev("o", 50*day), }, Tags: map[string]string{"demo": dg("k")}, }, } refs := []string{host + "/felis/paper:demo@" + dg("i")} got := targets(Plan(indexes, refs, host, now, day, 1)) want := []string{"felis/paper@j", "felis/paper@o", "felis/paper@r", "felis/paper@s"} if fmt.Sprint(got) != fmt.Sprint(want) { t.Fatalf("plan = %v, want %v", got, want) } } func TestParseRef(t *testing.T) { for _, c := range []struct { ref, repo, tag, digest string Loading Loading @@ -297,3 +331,73 @@ func TestRunThroughTheGate(t *testing.T) { t.Fatalf("upstream deletes = %q, want [%q]", deletes, want) } } // The gate reads an index's children off the blob store, so a pinned index keeps // its platform manifest end to end, and an unpinned one goes with it. func TestRunThroughTheGateKeepsAPinnedIndexWhole(t *testing.T) { root := t.TempDir() repoDir := filepath.Join(root, "docker/registry/v2/repositories/felis/paper/_manifests") hexOf := func(c string) string { return strings.Repeat(c, 64) } blob := func(hex, body string) { dir := filepath.Join(root, "docker/registry/v2/blobs/sha256", hex[:2], hex) if err := os.MkdirAll(dir, 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "data"), []byte(body), 0o644); err != nil { t.Fatal(err) } } for _, c := range []string{"1", "2", "3", "4"} { dir := filepath.Join(repoDir, "revisions/sha256", hexOf(c)) if err := os.MkdirAll(dir, 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "link"), []byte("sha256:"+hexOf(c)), 0o644); err != nil { t.Fatal(err) } } index := func(child string) string { return `{"schemaVersion":2,"mediaType":"application/vnd.oci.image.index.v1+json","manifests":[` + `{"mediaType":"application/vnd.oci.image.manifest.v1+json","digest":"sha256:` + hexOf(child) + `","size":2189,"platform":{"architecture":"arm64","os":"linux"}}]}` } blob(hexOf("1"), index("2")) // pinned index -> 2 blob(hexOf("2"), `{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","layers":[]}`) blob(hexOf("3"), index("4")) // unused index -> 4 var mu sync.Mutex var deleted []string up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { case r.URL.Path == "/v2/_catalog": _, _ = w.Write([]byte(`{"repositories":["felis/paper"]}`)) case r.Method == http.MethodDelete: mu.Lock() deleted = append(deleted, r.URL.Path[strings.LastIndex(r.URL.Path, ":")+1:][:1]) mu.Unlock() w.WriteHeader(http.StatusAccepted) default: http.NotFound(w, r) } })) t.Cleanup(up.Close) target, _ := url.Parse(up.URL) g := registrygate.New(target, map[string]string{registrygate.PrincipalPrune: "prune-secret"}, nil) g.DataDir = root gate := httptest.NewServer(g) t.Cleanup(gate.Close) refs := func(context.Context) ([]string, error) { return []string{host + "/felis/paper:demo@sha256:" + hexOf("1")}, nil } p := &Pruner{Registry: &Client{Endpoint: gate.URL, Token: "prune-secret"}, Host: host, Refs: refs, Now: func() time.Time { return time.Now().Add(48 * time.Hour) }} if _, err := p.Run(context.Background()); err != nil { t.Fatal(err) } mu.Lock() defer mu.Unlock() sort.Strings(deleted) if fmt.Sprint(deleted) != "[3 4]" { t.Fatalf("deleted = %v, want [3 4]: the pinned index 1 and its platform manifest 2 stay", deleted) } } Loading
internal/registrygate/index.go +43 −2 Changes for internal/registrygate/index.go: 43 added lines, 2 removed lines. Original line number Diff line number Diff line Loading @@ -3,6 +3,7 @@ package registrygate import ( "encoding/json" "errors" "io" "io/fs" "net/http" "os" Loading @@ -15,7 +16,7 @@ import ( // IndexPathPrefix serves the manifest index of one repository: // // GET /felis/manifests/<repo> {"revisions":[{"digest":…,"pushed":…}],"tags":{"<tag>":"<digest>"}} // GET /felis/manifests/<repo> {"revisions":[{"digest":…,"pushed":…,"children":[…]}],"tags":{"<tag>":"<digest>"}} // // The registry API can list a repository's tags but not its manifests, so a // manifest a tag moved off (every rebuild of :demo or :latest leaves one) is Loading @@ -39,6 +40,45 @@ type Index struct { type Revision struct { Digest string `json:"digest"` Pushed time.Time `json:"pushed"` // Children lists the manifests an image index (or Docker manifest list) names: // one per platform, plus BuildKit's attestation manifests. Each is a revision // of its own, untagged and never spelled in an image ref, yet a pull of the // index fetches them, so whoever keeps the index must keep them too. Empty for // a single-platform manifest. Children []string `json:"children,omitempty"` } // maxManifestBytes caps how much of a revision's blob is read to find its // children. It is the registry's own limit on a pushed manifest. const maxManifestBytes = 4 << 20 // manifestChildren reads a manifest blob off the filesystem driver's layout and // returns the digests it names when it is an index. A blob that is missing, // oversized or unparsable names nothing: a pull of it fails already, and the // pruner then treats it like any other revision. func manifestChildren(root, hex string) []string { f, err := os.Open(filepath.Join(root, "docker", "registry", "v2", "blobs", "sha256", hex[:2], hex, "data")) if err != nil { return nil } defer f.Close() // The manifests array is what makes an index, whatever mediaType says (an OCI // index may omit it); a single-platform manifest has layers and no manifests. var m struct { Manifests []struct { Digest string `json:"digest"` } `json:"manifests"` } if err := json.NewDecoder(io.LimitReader(f, maxManifestBytes)).Decode(&m); err != nil { return nil } var out []string for _, c := range m.Manifests { if digestRE.MatchString(c.Digest) { out = append(out, c.Digest) } } return out } // repoNameRE is the distribution reference grammar for a repository path. Every Loading Loading @@ -96,7 +136,8 @@ func ReadIndex(root, repo string) (*Index, error) { if err != nil { return nil, err } idx.Revisions = append(idx.Revisions, Revision{Digest: "sha256:" + e.Name(), Pushed: st.ModTime().UTC()}) idx.Revisions = append(idx.Revisions, Revision{Digest: "sha256:" + e.Name(), Pushed: st.ModTime().UTC(), Children: manifestChildren(root, e.Name())}) } sort.Slice(idx.Revisions, func(i, j int) bool { return idx.Revisions[i].Digest < idx.Revisions[j].Digest }) Loading
internal/registryprune/prune.go +29 −4 Changes for internal/registryprune/prune.go: 29 added lines, 4 removed lines. Original line number Diff line number Diff line Loading @@ -4,7 +4,8 @@ // previous manifest behind, a whitelist entry an admin removes keeps its image, // and each installer run adds a platform image. // // A manifest is kept when any of these hold: // A manifest is kept when any of these hold, and so is every manifest a kept // image index names (its per-platform images and attestations): // // - an image reference the platform still depends on names it: a whitelist // entry (by tag, by wildcard tag, or by digest), a server's spec (servers are Loading Loading @@ -211,16 +212,40 @@ func Plan(indexes map[string]*registrygate.Index, refs []string, host string, no } } for _, r := range idx.Revisions { t := Target{repo, r.Digest} if keep[t] || now.Sub(r.Pushed) < grace { continue if now.Sub(r.Pushed) < grace { keep[Target{repo, r.Digest}] = true } } keepChildren(repo, idx, keep) for _, r := range idx.Revisions { if t := (Target{repo, r.Digest}); !keep[t] { out = append(out, t) } } } return out } // keepChildren extends keep from each kept image index to the manifests it names: // the per-platform images and their attestations, which no ref spells but every // pull of the index fetches. An index can name another index, so it runs to a // fixed point. func keepChildren(repo string, idx *registrygate.Index, keep map[Target]bool) { for grew := true; grew; { grew = false for _, r := range idx.Revisions { if !keep[Target{repo, r.Digest}] { continue } for _, c := range r.Children { if t := (Target{repo, c}); !keep[t] { keep[t], grew = true, true } } } } } // newestTagged returns the n most recently pushed distinct digests any tag in idx // names. func newestTagged(idx *registrygate.Index, pushed map[string]time.Time, n int) []string { Loading
internal/registryprune/prune_test.go +104 −0 Changes for internal/registryprune/prune_test.go: 104 added lines, 0 removed lines. Original line number Diff line number Diff line Loading @@ -111,6 +111,40 @@ func TestPlanKeepsTheNewestTaggedPlatformImages(t *testing.T) { } } // A BuildKit push stores an image index plus one untagged revision per platform // and per attestation, and a server's spec pins the index. The children must stay // with it: the VM's test-one lost its arm64 manifest this way and could no longer // pull the paper image it was pinned to. func TestPlanKeepsTheManifestsAKeptIndexNames(t *testing.T) { withKids := func(r registrygate.Revision, kids ...string) registrygate.Revision { for _, k := range kids { r.Children = append(r.Children, dg(k)) } return r } indexes := map[string]*registrygate.Index{ "felis/paper": { Revisions: []registrygate.Revision{ // i: index pinned by a server, naming platform p and attestation q. withKids(rev("i", 30*day), "p", "q"), rev("p", 30*day), rev("q", 30*day), // j: an unused index; it goes with its children r and s. withKids(rev("j", 40*day), "r", "s"), rev("r", 40*day), rev("s", 40*day), // k: the newest tagged index names a nested index n, which names m. withKids(rev("k", 2*day), "n"), withKids(rev("n", 2*day), "m"), rev("m", 2*day), // o: a child of an index the registry no longer holds: goes. rev("o", 50*day), }, Tags: map[string]string{"demo": dg("k")}, }, } refs := []string{host + "/felis/paper:demo@" + dg("i")} got := targets(Plan(indexes, refs, host, now, day, 1)) want := []string{"felis/paper@j", "felis/paper@o", "felis/paper@r", "felis/paper@s"} if fmt.Sprint(got) != fmt.Sprint(want) { t.Fatalf("plan = %v, want %v", got, want) } } func TestParseRef(t *testing.T) { for _, c := range []struct { ref, repo, tag, digest string Loading Loading @@ -297,3 +331,73 @@ func TestRunThroughTheGate(t *testing.T) { t.Fatalf("upstream deletes = %q, want [%q]", deletes, want) } } // The gate reads an index's children off the blob store, so a pinned index keeps // its platform manifest end to end, and an unpinned one goes with it. func TestRunThroughTheGateKeepsAPinnedIndexWhole(t *testing.T) { root := t.TempDir() repoDir := filepath.Join(root, "docker/registry/v2/repositories/felis/paper/_manifests") hexOf := func(c string) string { return strings.Repeat(c, 64) } blob := func(hex, body string) { dir := filepath.Join(root, "docker/registry/v2/blobs/sha256", hex[:2], hex) if err := os.MkdirAll(dir, 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "data"), []byte(body), 0o644); err != nil { t.Fatal(err) } } for _, c := range []string{"1", "2", "3", "4"} { dir := filepath.Join(repoDir, "revisions/sha256", hexOf(c)) if err := os.MkdirAll(dir, 0o755); err != nil { t.Fatal(err) } if err := os.WriteFile(filepath.Join(dir, "link"), []byte("sha256:"+hexOf(c)), 0o644); err != nil { t.Fatal(err) } } index := func(child string) string { return `{"schemaVersion":2,"mediaType":"application/vnd.oci.image.index.v1+json","manifests":[` + `{"mediaType":"application/vnd.oci.image.manifest.v1+json","digest":"sha256:` + hexOf(child) + `","size":2189,"platform":{"architecture":"arm64","os":"linux"}}]}` } blob(hexOf("1"), index("2")) // pinned index -> 2 blob(hexOf("2"), `{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","layers":[]}`) blob(hexOf("3"), index("4")) // unused index -> 4 var mu sync.Mutex var deleted []string up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch { case r.URL.Path == "/v2/_catalog": _, _ = w.Write([]byte(`{"repositories":["felis/paper"]}`)) case r.Method == http.MethodDelete: mu.Lock() deleted = append(deleted, r.URL.Path[strings.LastIndex(r.URL.Path, ":")+1:][:1]) mu.Unlock() w.WriteHeader(http.StatusAccepted) default: http.NotFound(w, r) } })) t.Cleanup(up.Close) target, _ := url.Parse(up.URL) g := registrygate.New(target, map[string]string{registrygate.PrincipalPrune: "prune-secret"}, nil) g.DataDir = root gate := httptest.NewServer(g) t.Cleanup(gate.Close) refs := func(context.Context) ([]string, error) { return []string{host + "/felis/paper:demo@sha256:" + hexOf("1")}, nil } p := &Pruner{Registry: &Client{Endpoint: gate.URL, Token: "prune-secret"}, Host: host, Refs: refs, Now: func() time.Time { return time.Now().Add(48 * time.Hour) }} if _, err := p.Run(context.Background()); err != nil { t.Fatal(err) } mu.Lock() defer mu.Unlock() sort.Strings(deleted) if fmt.Sprint(deleted) != "[3 4]" { t.Fatalf("deleted = %v, want [3 4]: the pinned index 1 and its platform manifest 2 stay", deleted) } }