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.
This commit is contained in:
jochen
2026-10-05 18:13:23 +02:00
parent bb3cd6437b
commit 01c5ab2aab
13 changed files with 1252 additions and 32 deletions
+29 -7
View File
@@ -7,6 +7,8 @@ import (
"io"
"net/http"
"strings"
"github.com/novox/mesh-controller/internal/artifacts"
)
// Where built artifacts go.
@@ -66,22 +68,32 @@ func (r Registry) PublishImage(ctx context.Context, localTag, repository string)
return pinned, nil
}
// PublishArchive stores bytes as a blob and returns where to fetch them from.
// PublishArchive stores bytes as a blob, holds it by a manifest, and returns where to fetch them
// from.
//
// Two steps, which is the registry's own protocol: ask for somewhere to put it, then put it there
// naming the digest. The registry verifies the digest itself, so a blob that arrived corrupted is
// refused by the thing storing it rather than by the machine unpacking it a week later.
// Two steps for the blob, which is the registry's own protocol: ask for somewhere to put it, then
// put it there naming the digest. The registry verifies the digest itself, so a blob that arrived
// corrupted is refused by the thing storing it rather than by the machine unpacking it a week
// later.
//
// **Then a manifest that names it** (novox/hq issue 253, ADR 0189). The store's own collector
// marks only from manifests, and a blob no manifest names is collected however much the mesh
// means to keep it — so an archive published bare is an archive the first nightly collection
// deletes. The holder is composed from the digest and size alone (artifacts.Holder), which is
// what lets the sweep recompute it to backfill or let go without anything being recorded here.
// What a machine is told to fetch is the blob, exactly as before.
func (r Registry) PublishArchive(ctx context.Context, repository string, body []byte, digest string) (string, error) {
base := "http://" + r.Address + "/v2/" + repository
final := base + "/blobs/" + digest
// Already there. Blobs are immutable and named by their content, so this is not an
// optimisation — re-uploading would be asking the registry to store what it already has under
// the name it already has.
// the name it already has. It is still held: a blob published before holders existed is
// exactly the one a rebuild of the same source finds already there.
if there, err := r.has(ctx, final); err != nil {
return "", err
} else if there {
return final, nil
return final, r.hold(ctx, repository, digest, len(body))
}
start, err := http.NewRequestWithContext(ctx, http.MethodPost, base+"/blobs/uploads/", nil)
@@ -119,7 +131,17 @@ func (r Registry) PublishArchive(ctx context.Context, repository string, body []
said, _ := io.ReadAll(io.LimitReader(done.Body, 4096))
return "", fmt.Errorf("%s refused the blob: %s %s", base, done.Status, strings.TrimSpace(string(said)))
}
return final, nil
return final, r.hold(ctx, repository, digest, len(body))
}
// hold puts the manifest holding an archive beside it. A build whose archive could not be held is
// a failed build: recorded as published, it would be an archive the store's collector takes.
func (r Registry) hold(ctx context.Context, repository, digest string, size int) error {
store := artifacts.Store{Address: r.Address, HTTP: r.client()}
if _, err := store.HoldBlob(ctx, repository, digest, int64(size)); err != nil {
return fmt.Errorf("published %s/blobs/%s and could not hold it by a manifest: %w", repository, digest, err)
}
return nil
}
// has is whether this registry already holds what is at that URL.
+86 -6
View File
@@ -9,6 +9,8 @@ import (
"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.
@@ -18,9 +20,10 @@ import (
// there is not sent again, and that a tag is never accepted as a pin.
type fakeRegistry struct {
blobs map[string][]byte
uploads int
location string
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 {
@@ -28,6 +31,8 @@ func (f *fakeRegistry) serve(t *testing.T) *httptest.Server {
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) {
@@ -39,8 +44,28 @@ func (f *fakeRegistry) serve(t *testing.T) *httptest.Server {
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/"):
f.uploads++
where := f.location
if where == "" {
where = "/v2/upload/" + hex.EncodeToString([]byte("session"))
@@ -58,6 +83,7 @@ func (f *fakeRegistry) serve(t *testing.T) *httptest.Server {
return
}
f.blobs[digest] = body
f.puts[digest]++
w.WriteHeader(http.StatusCreated)
default:
w.WriteHeader(http.StatusNotFound)
@@ -92,6 +118,57 @@ func TestAnArchiveIsStoredAndFetchableByItsDigest(t *testing.T) {
}
}
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.
@@ -107,8 +184,11 @@ func TestABlobAlreadyThereIsNotSentAgain(t *testing.T) {
if _, err := r.PublishArchive(context.Background(), "shell/config", body, digest); err != nil {
t.Fatal(err)
}
if f.uploads != 1 {
t.Fatalf("the blob was uploaded %d times", f.uploads)
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))
}
}