Files
mesh-controller/internal/builder/registry_test.go
T
jochen 01c5ab2aab Hold every kept archive by a manifest so the store's collector keeps it (hq issue 253)
The store's garbage-collect marks only from manifests, and archives were
published as bare blobs, so the first real collection would delete every
archive the mesh keeps. PublishArchive now puts a deterministic OCI holder
manifest (empty config, one layer) beside each archive; the sweep holds every
kept archive before it lets anything go, which backfills existing bare blobs,
and lets go of an archive holder-first. A forgotten module no longer keeps its
five recent builds (ADR 0189). `collection [--json]` reports kept archives
held/unheld and what may be let go, so the dry run can be lifted on evidence.
2026-10-05 18:13:23 +02:00

255 lines
8.7 KiB
Go

package builder
import (
"context"
"crypto/sha256"
"encoding/hex"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/artifacts"
)
// An OCI registry as a content-addressed blob store, which is what it is.
//
// Against a stand-in rather than a real registry because what is under test is this side of the
// protocol — that the two steps happen in order, that the digest travels, that a blob already
// there is not sent again, and that a tag is never accepted as a pin.
type fakeRegistry struct {
blobs map[string][]byte
manifests map[string][]byte // "<repository>@<digest>"
puts map[string]int // blob uploads, by digest
location string
}
func (f *fakeRegistry) serve(t *testing.T) *httptest.Server {
t.Helper()
if f.blobs == nil {
f.blobs = map[string][]byte{}
}
f.manifests = map[string][]byte{}
f.puts = map[string]int{}
mux := http.NewServeMux()
server := httptest.NewServer(mux)
mux.HandleFunc("/v2/", func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/sha256:"):
digest := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]
if _, ok := f.blobs[digest]; ok {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
case strings.Contains(r.URL.Path, "/manifests/"):
// The archive's holder (novox/hq issue 253): asked for with an Accept, put by digest.
repository, digest, _ := strings.Cut(strings.TrimPrefix(r.URL.Path, "/v2/"), "/manifests/")
switch r.Method {
case http.MethodHead:
if _, ok := f.manifests[repository+"@"+digest]; ok &&
strings.Contains(r.Header.Get("Accept"), "application/vnd.oci.image.manifest.v1+json") {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodPut:
body, _ := io.ReadAll(r.Body)
sum := sha256.Sum256(body)
if digest != "sha256:"+hex.EncodeToString(sum[:]) {
w.WriteHeader(http.StatusBadRequest)
return
}
f.manifests[repository+"@"+digest] = body
w.WriteHeader(http.StatusCreated)
}
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/blobs/uploads/"):
where := f.location
if where == "" {
where = "/v2/upload/" + hex.EncodeToString([]byte("session"))
}
w.Header().Set("Location", where)
w.WriteHeader(http.StatusAccepted)
case r.Method == http.MethodPut:
body, _ := io.ReadAll(r.Body)
digest := r.URL.Query().Get("digest")
sum := sha256.Sum256(body)
if digest != "sha256:"+hex.EncodeToString(sum[:]) {
// A registry verifies what it is given, which is why a corrupt blob is refused
// here rather than by a machine unpacking it a week later.
w.WriteHeader(http.StatusBadRequest)
return
}
f.blobs[digest] = body
f.puts[digest]++
w.WriteHeader(http.StatusCreated)
default:
w.WriteHeader(http.StatusNotFound)
}
})
t.Cleanup(server.Close)
return server
}
func registryFor(t *testing.T, f *fakeRegistry) Registry {
t.Helper()
server := f.serve(t)
return Registry{Address: strings.TrimPrefix(server.URL, "http://")}
}
func TestAnArchiveIsStoredAndFetchableByItsDigest(t *testing.T) {
f := &fakeRegistry{}
r := registryFor(t, f)
body := []byte("a theme")
sum := sha256.Sum256(body)
digest := "sha256:" + hex.EncodeToString(sum[:])
where, err := r.PublishArchive(context.Background(), "shell/config", body, digest)
if err != nil {
t.Fatal(err)
}
if !strings.HasSuffix(where, "/blobs/"+digest) {
t.Fatalf("what a machine is told to fetch is %q, which does not name the digest", where)
}
if string(f.blobs[digest]) != "a theme" {
t.Fatalf("the registry holds %q", f.blobs[digest])
}
}
func TestAnArchiveIsPublishedWithAManifestHoldingIt(t *testing.T) {
// The store's collector marks only from manifests, so a bare blob is one the first collection
// deletes, kept or not (novox/hq issue 253, ADR 0189). Every archive goes out held, by the
// holder the sweep can compute for itself from the digest and size.
f := &fakeRegistry{}
r := registryFor(t, f)
body := []byte("a theme")
sum := sha256.Sum256(body)
digest := "sha256:" + hex.EncodeToString(sum[:])
where, err := r.PublishArchive(context.Background(), "shell/config", body, digest)
if err != nil {
t.Fatal(err)
}
if !strings.HasSuffix(where, "/v2/shell/config/blobs/"+digest) {
t.Fatalf("a machine is told to fetch %q; the blob is still what is fetched", where)
}
manifest, holder := artifacts.Holder(digest, int64(len(body)))
if got := f.manifests["shell/config@"+holder]; string(got) != string(manifest) {
t.Fatalf("the store holds %v; want the holder %s in the archive's own repository", f.manifests, holder)
}
if !strings.Contains(string(manifest), `"digest":"`+digest+`"`) {
t.Fatalf("the holder does not name the archive: %s", manifest)
}
empty := sha256.Sum256([]byte("{}"))
if _, ok := f.blobs["sha256:"+hex.EncodeToString(empty[:])]; !ok {
t.Fatal("the holder's empty config was never put, and a registry refuses a manifest without it")
}
}
func TestAnArchiveAlreadyThereIsStillHeld(t *testing.T) {
// A rebuild of the same source finds its archive already there — often one published bare,
// before holders. It is not sent again, and it is held.
f := &fakeRegistry{}
r := registryFor(t, f)
body := []byte("published bare")
sum := sha256.Sum256(body)
digest := "sha256:" + hex.EncodeToString(sum[:])
f.blobs[digest] = body
if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil {
t.Fatal(err)
}
if f.puts[digest] != 0 {
t.Fatal("a blob already there was sent again")
}
if _, holder := artifacts.Holder(digest, int64(len(body))); f.manifests["shell/config@"+holder] == nil {
t.Fatal("a blob already there was left bare")
}
}
func TestABlobAlreadyThereIsNotSentAgain(t *testing.T) {
// Not an optimisation: blobs are named by their content, so re-uploading is asking the
// registry to store what it already has under the name it already has.
f := &fakeRegistry{}
r := registryFor(t, f)
body := []byte("a theme")
sum := sha256.Sum256(body)
digest := "sha256:" + hex.EncodeToString(sum[:])
if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil {
t.Fatal(err)
}
if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil {
t.Fatal(err)
}
if f.puts[digest] != 1 {
t.Fatalf("the blob was uploaded %d times", f.puts[digest])
}
if len(f.manifests) != 1 {
t.Fatalf("publishing twice left %d holders; want the one", len(f.manifests))
}
}
func TestAnAbsoluteUploadLocationIsFollowed(t *testing.T) {
// A registry may answer with a full URL or with a path. Both happen in the wild, and a client
// that handles one produces a confusing failure against the other.
f := &fakeRegistry{}
server := f.serve(t)
f.location = server.URL + "/v2/upload/session?state=abc"
r := Registry{Address: strings.TrimPrefix(server.URL, "http://")}
body := []byte("x")
sum := sha256.Sum256(body)
digest := "sha256:" + hex.EncodeToString(sum[:])
if _, err := r.PublishArchive(context.Background(), "a/b", body, digest); err != nil {
t.Fatal(err)
}
if _, ok := f.blobs[digest]; !ok {
t.Fatal("the blob did not arrive when the location carried a query")
}
}
func TestATagIsNeverAcceptedAsAPin(t *testing.T) {
// A tag can be made to point at something else, and a bundle is applied on machines with no
// mesh to ask about anything.
r := Registry{Address: "registry.invalid", Run: func(
_ context.Context, _, name string, args ...string) (string, error) {
if name == "docker" && len(args) > 0 && args[0] == "inspect" {
return "registry.invalid/shell/server:latest\n", nil
}
return "", nil
}}
_, err := r.PublishImage(context.Background(), "local", "shell/server")
if err == nil {
t.Fatal("an image referred to by tag was accepted")
}
if !strings.Contains(err.Error(), "pinned by digest") {
t.Fatalf("refused for the wrong reason: %v", err)
}
}
func TestAPushedImageComesBackPinned(t *testing.T) {
pinned := "registry.invalid/shell/server@sha256:" + strings.Repeat("a", 64)
var ran []string
r := Registry{Address: "registry.invalid", Run: func(
_ context.Context, _, name string, args ...string) (string, error) {
ran = append(ran, name+" "+strings.Join(args, " "))
if name == "docker" && len(args) > 0 && args[0] == "inspect" {
return pinned + "\n", nil
}
return "", nil
}}
got, err := r.PublishImage(context.Background(), "local", "shell/server")
if err != nil {
t.Fatal(err)
}
if got != pinned {
t.Fatalf("got %q", got)
}
if len(ran) < 2 || !strings.HasPrefix(ran[0], "docker tag") || !strings.HasPrefix(ran[1], "docker push") {
t.Fatalf("it did not tag and then push: %v", ran)
}
}