diff --git a/cmd/mesh-controller/collect.go b/cmd/mesh-controller/collect.go index 6d01e2a..d30344e 100644 --- a/cmd/mesh-controller/collect.go +++ b/cmd/mesh-controller/collect.go @@ -31,7 +31,12 @@ func collect(ctx context.Context, inv *inventory.Inventory) { fmt.Fprintf(os.Stderr, "could not work out what the artifact store may let go of: %v\n", err) return } - if len(references) == 0 { + kept, err := inv.KeptArchives(ctx) + if err != nil { + fmt.Fprintf(os.Stderr, "could not work out which archives the artifact store keeps: %v\n", err) + return + } + if len(references) == 0 && len(kept) == 0 { return } shelf, err := inv.Catalogue(ctx) @@ -59,6 +64,25 @@ func collect(ctx context.Context, inv *inventory.Inventory) { defer stop() store := artifacts.Store{Address: address} + // **Hold before letting go** (novox/hq issue 253). The store's collector keeps only what a + // manifest names, and archives were published as bare blobs, so every kept archive is first + // held by its manifest — which backfills the ones published before holders, a few at a time + // as builds come, and is two HEADs each once done. A kept archive that could not be held stops + // the sweep before it deletes anything: "everything kept is held" is the precondition the + // collector's safety rests on, and a store refusing a hold would refuse the deletes too. + wrote, missing, err := holdKept(within, store, kept) + if wrote > 0 { + fmt.Fprintf(os.Stderr, "the artifact store now holds %d more kept archive(s) by a manifest\n", wrote) + } + if missing > 0 { + fmt.Fprintf(os.Stderr, "%d archive(s) the mesh keeps are not in the artifact store at all; "+ + "`collection` lists them\n", missing) + } + if err != nil { + fmt.Fprintf(os.Stderr, "not every kept archive could be held, so nothing was let go: %v\n", err) + return + } + var done []string var left, skipped int for i, reference := range references { @@ -114,6 +138,34 @@ func collect(ctx context.Context, inv *inventory.Inventory) { } } +// holdKept holds every kept archive by its manifest, stopping at the first refusal by the store. +// Answers how many holders it wrote and how many kept archives the store does not have. +// +// A missing archive is counted rather than fatal: there is nothing to hold, and that is a fact +// for an operator to read (`collection`), not a reason to stop collecting what is not kept. A +// reference the store cannot be asked about is skipped as the deletion loop skips one +// (novox/hq issue 226). +func holdKept(ctx context.Context, store artifacts.Store, kept []string) (wrote, missing int, err error) { + for _, reference := range kept { + if err := ctx.Err(); err != nil { + return wrote, missing, fmt.Errorf("ran out of time before %s: %w", reference, err) + } + did, err := store.Hold(ctx, reference) + switch { + case err == nil: + if did { + wrote++ + } + case errors.Is(err, artifacts.Gone): + missing++ + case errors.Is(err, artifacts.ErrNotOurs): + default: + return wrote, missing, fmt.Errorf("holding %s: %w", reference, err) + } + } + return wrote, missing, nil +} + // mostPerSweep is how many artifacts one sweep will ask about. Enough that a mesh building // several times a day converges within days of this landing; small enough that no single build // waits on the whole backlog. diff --git a/cmd/mesh-controller/collection.go b/cmd/mesh-controller/collection.go new file mode 100644 index 0000000..6e2520d --- /dev/null +++ b/cmd/mesh-controller/collection.go @@ -0,0 +1,166 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "os" + "strings" + + "github.com/novox/mesh-controller/internal/artifacts" +) + +// What the artifact store keeps, whether each kept archive is held, and what the sweep may let go +// (novox/hq issue 253, ADR 0189). +// +// **The question to answer before the store's collector runs for real.** The collector deletes +// every blob no manifest names, and archives were published as bare blobs, so the nightly step +// runs `--dry-run` until every archive the mesh keeps is held by its manifest. The sweep holds +// them as builds come; this says how far that has got — "0 unheld" is the number that lets the +// dry run go. +// +// Reads and changes nothing: each kept archive is asked about with HEADs only. Reached over the +// console through the mesh-controller seat's `command` verb (`collection --json`), which needs no +// new verb in the seat's row. + +type collectionReport struct { + // Store is the artifact store as this machine reached it; empty when it is not on the network. + Store string `json:"store"` + // KeptArchives is how many archives the mesh keeps, for either reason. + KeptArchives int `json:"kept_archives"` + // Held is how many of them the store holds by their manifest. + Held int `json:"held"` + // Unheld are the kept archives the collector would delete tonight if it ran for real. + Unheld []string `json:"unheld"` + // Missing are kept archives the store does not have at all. + Missing []string `json:"missing"` + // Unasked is how many could not be asked about, and why the asking stopped. + Unasked int `json:"unasked"` + Stopped string `json:"stopped,omitempty"` + // Eligible is what the sweep may let go of: made by the mesh, kept for no reason, not yet + // collected — split by kind. + Eligible int `json:"eligible"` + EligibleImages int `json:"eligible_images"` + EligibleArchives int `json:"eligible_archives"` + // SafeToCollect is whether every kept archive was asked about and every one is held. + SafeToCollect bool `json:"safe_to_collect"` +} + +func collectionCommand(ctx context.Context, args []string) error { + set := flag.NewFlagSet("collection", flag.ContinueOnError) + asJSON := set.Bool("json", false, "answer as JSON") + positionals, err := parseAround(set, args) + if err != nil { + return err + } + if len(positionals) != 0 { + return errors.New("collection [--json]") + } + + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + inv := open.inventory + + kept, err := inv.KeptArchives(ctx) + if err != nil { + return err + } + eligible, err := inv.ToCollect(ctx) + if err != nil { + return err + } + report := collectionReport{KeptArchives: len(kept), Eligible: len(eligible), Unheld: []string{}, Missing: []string{}} + for _, reference := range eligible { + if strings.Contains(reference, "/blobs/") { + report.EligibleArchives++ + } else { + report.EligibleImages++ + } + } + + shelf, err := inv.Catalogue(ctx) + if err != nil { + return err + } + report.Store, err = artifactStoreAddress(ctx, inv, shelf, "") + if err != nil { + return err + } + if report.Store == "" { + report.Unasked = len(kept) + report.Stopped = "this mesh has no artifact store on its network" + } else { + report.Unasked, report.Stopped = askHeld(ctx, artifacts.Store{Address: report.Store}, kept, &report) + } + report.SafeToCollect = report.Unasked == 0 && len(report.Unheld) == 0 + + if *asJSON { + encoder := json.NewEncoder(os.Stdout) + encoder.SetIndent("", " ") + return encoder.Encode(report) + } + printCollection(report) + return nil +} + +// askHeld asks the store about each kept archive, stopping at the first answer that is not about +// the archive: a store that cannot be reached for one cannot be for the next, and a page of +// identical failures says less than one line. +func askHeld(ctx context.Context, store artifacts.Store, kept []string, report *collectionReport) (int, string) { + for i, reference := range kept { + held, err := store.Held(ctx, reference) + switch { + case err == nil && held: + report.Held++ + case err == nil: + report.Unheld = append(report.Unheld, reference) + case errors.Is(err, artifacts.Gone): + report.Missing = append(report.Missing, reference) + case errors.Is(err, artifacts.ErrNotOurs): + // KeptArchives names only the mesh's own; counted as unasked if one ever is not. + report.Unasked++ + default: + return report.Unasked + len(kept) - i, fmt.Sprintf("asking about %s: %v", reference, err) + } + } + return report.Unasked, "" +} + +func printCollection(r collectionReport) { + store := r.Store + if store == "" { + store = "(not on the network)" + } + fmt.Printf("artifact store %s\n", store) + fmt.Printf("kept archives %d\n", r.KeptArchives) + fmt.Printf(" held %d\n", r.Held) + fmt.Printf(" unheld %d\n", len(r.Unheld)) + fmt.Printf(" missing %d\n", len(r.Missing)) + if r.Unasked > 0 { + fmt.Printf(" not asked %d (%s)\n", r.Unasked, r.Stopped) + } + fmt.Printf("eligible to let go %d (%d images, %d archives)\n", r.Eligible, r.EligibleImages, r.EligibleArchives) + if len(r.Unheld) > 0 { + fmt.Println("\nunheld — the store's collector would delete these; the next build's sweep holds them:") + for _, reference := range r.Unheld { + fmt.Printf(" %s\n", reference) + } + } + if len(r.Missing) > 0 { + fmt.Println("\nmissing — kept by the mesh, not in the store:") + for _, reference := range r.Missing { + fmt.Printf(" %s\n", reference) + } + } + fmt.Println() + if r.SafeToCollect { + fmt.Println("every kept archive is held: the store's collector may run for real") + } else { + fmt.Println("NOT every kept archive is known to be held: keep the store's collector on --dry-run") + } +} diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index 0f2561f..b15f40d 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -73,6 +73,8 @@ func run() error { return askCommand(ctx, args[1:]) case "builds": return buildsCommand(ctx, args[1:]) + case "collection": + return collectionCommand(ctx, args[1:]) case "plans": return plansCommand(ctx, args[1:]) case "pin": @@ -200,6 +202,7 @@ func usage() { build --behind build every module the mesh holds older than its source build --on rebuild every module that stands on this module's artifacts, bases first builds [] what has been built lately, and what came of it + collection [--json] kept archives held/unheld by a manifest, and what the sweep may let go builder issue a broker account for a build machine, scoped to build work, delivered as the builder module's broker secret (module add it first) licence add|list|use|key model access, under the name a person calls it diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 9d2bb65..3591da5 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -198,7 +198,7 @@ func argvFor(verb string, args map[string]any) ([]string, error) { } // jsonVerbs are the verbs whose command speaks JSON, so the answer carries it as data as well. -var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true} +var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true, "collection": true} // runVerb runs this binary with the given command line and gathers what it said. func runVerb(ctx context.Context, argv []string) (verbAnswer, error) { diff --git a/internal/artifacts/hold.go b/internal/artifacts/hold.go new file mode 100644 index 0000000..d8d1d9b --- /dev/null +++ b/internal/artifacts/hold.go @@ -0,0 +1,345 @@ +package artifacts + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "strconv" + "strings" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// Every archive the store keeps is held by a manifest (novox/hq issue 253, ADR 0189). +// +// **The store's collector marks only from manifests.** The mesh's store is a stock registry, and +// its nightly `registry garbage-collect` walks every manifest in every repository, marks the blobs +// those manifests name, and deletes every blob it did not mark. An image is a manifest, so what +// the mesh keeps of an image survives. An archive was not: the builder put it in the store as a +// bare blob — upload, then `PUT ?digest=` — and nothing in the store names it. To the collector a +// bare blob is unreferenced, so the first real collection would have deleted every archive the +// mesh holds, kept or not, and every machine pinning a bundle would have found it gone. The +// collector runs `--dry-run` until this is true. +// +// **So each archive gets a holder**: the smallest OCI image manifest that names it — the empty +// config, one layer, nothing else — put in the archive's own repository, by digest, untagged. The +// collector marks it and so keeps the archive; the sweep lets go of an archive by deleting its +// holder first, which is what lets the bytes go at the next collection. +// +// **Nothing a machine reads changes.** The recorded reference stays +// `artifact-store:////blobs/sha256:…`, and machines fetch the blob exactly as +// before. The holder is the store's bookkeeping, not a second way to reach anything. +// +// **Deterministic, so it never needs recording.** The holder is composed from the archive's digest +// and size alone, in a fixed field order with no timestamps or annotations, so the sweep can +// compute which manifest holds any archive from the reference it already has plus one HEAD for the +// size. No schema change, no second record that could disagree with the store. + +const ( + // mediaManifest is the type a holder is put and asked for as. + mediaManifest = "application/vnd.oci.image.manifest.v1+json" + // mediaEmpty is the OCI empty descriptor's type: a config that says nothing, for a manifest + // whose only purpose is to name its layer. + mediaEmpty = "application/vnd.oci.empty.v1+json" + // emptyDigest is the digest of `{}`, the empty config's content, fixed by the OCI spec. + emptyDigest = "sha256:44136fa355b3678a1146ad16f7e8649e94fb4fc21fe77e8310c060f61caaff8a" + // mediaArchive is the layer type an archive is held as. Every archive the builder publishes is + // `pack`'s gzipped tar, so this is the true type and not a placeholder — and it is a constant, + // not read from anywhere, because the holder must be recomputable from the reference alone. + mediaArchive = "application/vnd.oci.image.layer.v1.tar+gzip" +) + +// emptyConfig is the content emptyDigest names. +var emptyConfig = []byte("{}") + +// manifestAccept is what a manifest is asked for as. A registry answers a manifest HEAD only in a +// type the caller named, and answers 404 to a bare one for a manifest it holds perfectly well +// (measured 2026-09-28; internal/builder/registry.go says how that was found). +var manifestAccept = []string{ + mediaManifest, + "application/vnd.docker.distribution.manifest.v2+json", +} + +type descriptor struct { + MediaType string `json:"mediaType"` + Digest string `json:"digest"` + Size int64 `json:"size"` +} + +type holderManifest struct { + SchemaVersion int `json:"schemaVersion"` + MediaType string `json:"mediaType"` + Config descriptor `json:"config"` + Layers []descriptor `json:"layers"` +} + +// Holder is the manifest that holds an archive in the store, and its digest. +// +// A pure function of the archive's digest and size: the same two in give the same bytes out, +// always, because `encoding/json` writes a struct's fields in their declared order and there is +// nothing here that varies by when or where it was composed. +func Holder(digest string, size int64) (body []byte, holder string) { + body, err := json.Marshal(holderManifest{ + SchemaVersion: 2, + MediaType: mediaManifest, + Config: descriptor{MediaType: mediaEmpty, Digest: emptyDigest, Size: int64(len(emptyConfig))}, + Layers: []descriptor{{MediaType: mediaArchive, Digest: digest, Size: size}}, + }) + if err != nil { + // Marshalling a struct of strings and integers cannot fail. + panic(err) + } + sum := sha256.Sum256(body) + return body, "sha256:" + hex.EncodeToString(sum[:]) +} + +// Hold makes sure the store holds this archive by a manifest, and says whether it had to write one. +// +// Takes a reference as the mesh records it. An image is its own manifest and needs no holder, so +// it answers false and nothing is asked. Idempotent: a holder already there is left alone, which +// is what lets the sweep run it over every kept archive on every build and so backfill the bare +// blobs published before holders existed (novox/hq issue 253). +// +// Gone when the store does not hold the archive at all: there is nothing to hold, and that is a +// fact the caller reports rather than one this invents a remedy for. +func (s Store) Hold(ctx context.Context, reference string) (bool, error) { + repository, digest, archive, err := s.archive(reference) + if err != nil || !archive { + return false, err + } + size, err := s.blobSize(ctx, repository, digest) + if err != nil { + return false, err + } + return s.HoldBlob(ctx, repository, digest, size) +} + +// Held is whether the store holds this archive by its manifest. Asks and changes nothing — the +// question an operator needs answered with "none unheld" before the collector is let loose. +// +// An image answers true: it is its own manifest. An archive the store does not have answers Gone. +func (s Store) Held(ctx context.Context, reference string) (bool, error) { + repository, digest, archive, err := s.archive(reference) + if err != nil { + return false, err + } + if !archive { + return true, nil + } + size, err := s.blobSize(ctx, repository, digest) + if err != nil { + return false, err + } + _, holder := Holder(digest, size) + return s.has(ctx, s.url(repository, "manifests", holder), manifestAccept...) +} + +// HoldBlob puts the holder for a blob of this digest and size into its repository, unless it is +// there already. Answers whether it wrote one. +// +// The builder calls this with the size it has just uploaded; the sweep, through Hold, with the size +// the store reports. Both arrive at the same holder, which is the point of composing it. +func (s Store) HoldBlob(ctx context.Context, repository, digest string, size int64) (bool, error) { + if s.Address == "" { + return false, fmt.Errorf("this mesh has no artifact store on its network to hold %s/%s in", repository, digest) + } + body, holder := Holder(digest, size) + there, err := s.has(ctx, s.url(repository, "manifests", holder), manifestAccept...) + if err != nil { + return false, err + } + if there { + return false, nil + } + + // The config must be in the repository before a manifest naming it is accepted: a registry + // refuses a manifest whose blobs it cannot find there, which is the property that makes a + // holder mean something. + if err := s.putBlob(ctx, repository, emptyDigest, emptyConfig); err != nil { + return false, err + } + + // **By digest, never by tag.** A tag would be one more name to move and one more thing the + // collector's `--delete-untagged` would read as meaningful; the mesh names nothing by tag that + // it pins by digest, and an untagged manifest is kept by plain collection. + request, err := http.NewRequestWithContext(ctx, http.MethodPut, + s.url(repository, "manifests", holder), bytes.NewReader(body)) + if err != nil { + return false, err + } + request.Header.Set("Content-Type", mediaManifest) + response, err := s.client().Do(request) + if err != nil { + return false, err + } + defer response.Body.Close() + if response.StatusCode != http.StatusCreated { + said, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + return false, fmt.Errorf("the artifact store refused to hold %s/%s: %s %s", + repository, digest, response.Status, strings.TrimSpace(string(said))) + } + return true, nil +} + +// letGoOfHolder deletes the manifest holding an archive, before the archive's own link goes. +// +// **Holder first.** Deleting the blob link alone leaves a manifest still naming the blob, and the +// collector would keep its bytes for ever on the strength of it — the sweep would record the +// archive collected while the disk said otherwise. Deleting the holder first and failing before +// the link goes leaves an unheld archive that the next sweep still offers, which is safe. +// +// A store that no longer has the blob answers Gone: without its size the holder cannot be named, +// and without the blob there is nothing left for a holder to keep. A store that never had a holder +// for it — an archive published before holders, never backfilled — answers 404 to the delete, and +// that is the outcome wanted. +func (s Store) letGoOfHolder(ctx context.Context, repository, digest string) error { + size, err := s.blobSize(ctx, repository, digest) + if err != nil { + return err + } + _, holder := Holder(digest, size) + err = s.remove(ctx, s.url(repository, "manifests", holder), repository+"/manifests/"+holder) + if err == Gone { + return nil + } + return err +} + +// archive reads a recorded reference into its repository and digest, and whether it is an archive +// at all. Refuses as ErrNotOurs anything the mesh did not put in its own store. +func (s Store) archive(reference string) (repository, digest string, archive bool, err error) { + path, kept := catalogue.InArtifactStore(reference) + if !kept { + return "", "", false, fmt.Errorf("%w: %s", ErrNotOurs, reference) + } + if s.Address == "" { + return "", "", false, fmt.Errorf("this mesh has no artifact store on its network to ask about %s", reference) + } + repository, kind, digest, err := split(path) + if err != nil { + return "", "", false, err + } + return repository, digest, kind == "blobs", nil +} + +// blobSize is how large the store says a blob is; Gone when it does not have it. +func (s Store) blobSize(ctx context.Context, repository, digest string) (int64, error) { + request, err := http.NewRequestWithContext(ctx, http.MethodHead, s.url(repository, "blobs", digest), nil) + if err != nil { + return 0, err + } + response, err := s.client().Do(request) + if err != nil { + return 0, fmt.Errorf("cannot reach the artifact store at %s: %w", s.Address, err) + } + defer response.Body.Close() + switch response.StatusCode { + case http.StatusOK: + case http.StatusNotFound: + return 0, Gone + default: + return 0, fmt.Errorf("the artifact store answered %s for %s/blobs/%s", response.Status, repository, digest) + } + // Read from the header rather than ContentLength: a HEAD's ContentLength is what the response + // says it would have sent, which Go reports faithfully, but a proxy in between is free to drop + // it, and the header is what the registry itself wrote. + if length := response.Header.Get("Content-Length"); length != "" { + if n, err := strconv.ParseInt(length, 10, 64); err == nil && n >= 0 { + return n, nil + } + } + if response.ContentLength >= 0 { + return response.ContentLength, nil + } + return 0, fmt.Errorf("the artifact store holds %s/blobs/%s and will not say how large it is", repository, digest) +} + +// putBlob uploads a small blob unless the repository already has it: ask where, then put it there +// naming the digest — the registry's own two steps, the same the builder takes for an archive. +func (s Store) putBlob(ctx context.Context, repository, digest string, body []byte) error { + if there, err := s.has(ctx, s.url(repository, "blobs", digest)); err != nil { + return err + } else if there { + return nil + } + start, err := http.NewRequestWithContext(ctx, http.MethodPost, + "http://"+s.Address+"/v2/"+repository+"/blobs/uploads/", nil) + if err != nil { + return err + } + begun, err := s.client().Do(start) + if err != nil { + return fmt.Errorf("cannot start an upload to %s: %w", repository, err) + } + begun.Body.Close() + if begun.StatusCode != http.StatusAccepted { + return fmt.Errorf("the artifact store answered %s when asked where to put a blob in %s", begun.Status, repository) + } + where := begun.Header.Get("Location") + if where == "" { + return fmt.Errorf("the artifact store accepted an upload to %s and said nowhere to put it", repository) + } + if strings.HasPrefix(where, "/") { + where = "http://" + s.Address + where + } + separator := "?" + if strings.Contains(where, "?") { + separator = "&" + } + put, err := http.NewRequestWithContext(ctx, http.MethodPut, where+separator+"digest="+digest, bytes.NewReader(body)) + if err != nil { + return err + } + put.Header.Set("Content-Type", "application/octet-stream") + done, err := s.client().Do(put) + if err != nil { + return err + } + defer done.Body.Close() + if done.StatusCode != http.StatusCreated { + said, _ := io.ReadAll(io.LimitReader(done.Body, 4096)) + return fmt.Errorf("the artifact store refused a blob in %s: %s %s", repository, done.Status, strings.TrimSpace(string(said))) + } + return nil +} + +// has is whether the store answers 200 for a HEAD at that URL. +func (s Store) has(ctx context.Context, url string, accept ...string) (bool, error) { + request, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil) + if err != nil { + return false, err + } + for _, media := range accept { + request.Header.Add("Accept", media) + } + response, err := s.client().Do(request) + if err != nil { + return false, fmt.Errorf("cannot reach the artifact store at %s: %w", s.Address, err) + } + defer response.Body.Close() + switch response.StatusCode { + case http.StatusOK: + return true, nil + case http.StatusNotFound: + return false, nil + default: + return false, fmt.Errorf("the artifact store answered %s for %s", response.Status, url) + } +} + +func (s Store) url(repository, kind, digest string) string { + return "http://" + s.Address + "/v2/" + repository + "/" + kind + "/" + digest +} + +func (s Store) client() *http.Client { + if s.HTTP != nil { + return s.HTTP + } + return &http.Client{Timeout: 30 * time.Second} +} diff --git a/internal/artifacts/hold_test.go b/internal/artifacts/hold_test.go new file mode 100644 index 0000000..8f241f6 --- /dev/null +++ b/internal/artifacts/hold_test.go @@ -0,0 +1,268 @@ +package artifacts + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/http/httptest" + "strconv" + "strings" + "sync" + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// Every kept archive is held by a manifest (novox/hq issue 253, ADR 0189). +// +// Against an in-memory registry that keeps blobs and manifests per repository and refuses what a +// registry refuses — a blob whose digest does not match, a manifest whose digest does not match +// or whose blobs the repository does not have, a manifest asked for without an Accept naming its +// type. What is asserted is this side's decisions; the live test below asserts the registry's. + +type memRegistry struct { + mu sync.Mutex + blobs map[string][]byte // repository + "@" + digest + manifests map[string][]byte // repository + "@" + digest + writes []string // every PUT and DELETE, as "METHOD path" +} + +func digestOf(body []byte) string { + sum := sha256.Sum256(body) + return "sha256:" + hex.EncodeToString(sum[:]) +} + +func (m *memRegistry) serve(t *testing.T) Store { + t.Helper() + m.blobs = map[string][]byte{} + m.manifests = map[string][]byte{} + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + m.mu.Lock() + defer m.mu.Unlock() + path := strings.TrimPrefix(r.URL.Path, "/v2/") + if r.Method == http.MethodPut || r.Method == http.MethodDelete { + m.writes = append(m.writes, r.Method+" "+r.URL.Path) + } + switch { + case r.Method == http.MethodPost && strings.HasSuffix(path, "/blobs/uploads/"): + repository := strings.TrimSuffix(path, "/blobs/uploads/") + w.Header().Set("Location", "/upload/"+repository+"?state=x") + w.WriteHeader(http.StatusAccepted) + case r.Method == http.MethodPut && strings.HasPrefix(r.URL.Path, "/upload/"): + repository := strings.TrimPrefix(r.URL.Path, "/upload/") + body, _ := io.ReadAll(r.Body) + digest := r.URL.Query().Get("digest") + if digest != digestOf(body) { + w.WriteHeader(http.StatusBadRequest) + return + } + m.blobs[repository+"@"+digest] = body + w.WriteHeader(http.StatusCreated) + case strings.Contains(path, "/blobs/"): + repository, digest, _ := strings.Cut(path, "/blobs/") + key := repository + "@" + digest + body, ok := m.blobs[key] + if !ok { + w.WriteHeader(http.StatusNotFound) + return + } + switch r.Method { + case http.MethodHead: + w.Header().Set("Content-Length", strconv.Itoa(len(body))) + w.WriteHeader(http.StatusOK) + case http.MethodDelete: + delete(m.blobs, key) + w.WriteHeader(http.StatusAccepted) + default: + w.WriteHeader(http.StatusMethodNotAllowed) + } + case strings.Contains(path, "/manifests/"): + repository, digest, _ := strings.Cut(path, "/manifests/") + key := repository + "@" + digest + switch r.Method { + case http.MethodHead: + if _, ok := m.manifests[key]; !ok || !strings.Contains(r.Header.Get("Accept"), mediaManifest) { + w.WriteHeader(http.StatusNotFound) + return + } + w.WriteHeader(http.StatusOK) + case http.MethodPut: + body, _ := io.ReadAll(r.Body) + if digest != digestOf(body) { + w.WriteHeader(http.StatusBadRequest) + return + } + var named holderManifest + if err := json.Unmarshal(body, &named); err != nil { + w.WriteHeader(http.StatusBadRequest) + return + } + for _, d := range append([]descriptor{named.Config}, named.Layers...) { + if _, ok := m.blobs[repository+"@"+d.Digest]; !ok { + w.WriteHeader(http.StatusBadRequest) + fmt.Fprintf(w, "MANIFEST_BLOB_UNKNOWN %s", d.Digest) + return + } + } + m.manifests[key] = body + w.WriteHeader(http.StatusCreated) + case http.MethodDelete: + if _, ok := m.manifests[key]; !ok { + w.WriteHeader(http.StatusNotFound) + return + } + delete(m.manifests, key) + w.WriteHeader(http.StatusAccepted) + } + default: + w.WriteHeader(http.StatusNotFound) + } + })) + t.Cleanup(server.Close) + return Store{Address: strings.TrimPrefix(server.URL, "http://")} +} + +// bare puts an archive in the store the way the builder did before holders: a blob, nothing more. +func (m *memRegistry) bare(repository string, body []byte) string { + m.mu.Lock() + defer m.mu.Unlock() + digest := digestOf(body) + m.blobs[repository+"@"+digest] = body + return catalogue.ArtifactStoreScheme + repository + "/blobs/" + digest +} + +func TestTheHolderIsComposedFromTheDigestAndSizeAlone(t *testing.T) { + // The sweep must arrive at the very manifest the builder wrote, with nothing recorded between + // them. Same inputs, same bytes — and a different size is a different holder, so a holder can + // never be mistaken for one of a different blob. + digest := "sha256:" + strings.Repeat("a", 64) + one, first := Holder(digest, 42) + two, second := Holder(digest, 42) + if !bytes.Equal(one, two) || first != second { + t.Fatalf("the same archive composed two holders:\n%s\n%s", one, two) + } + if _, other := Holder(digest, 43); other == first { + t.Fatal("a different size composed the same holder") + } + want := `{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json",` + + `"config":{"mediaType":"application/vnd.oci.empty.v1+json",` + + `"digest":"sha256:44136fa355b3678a1146ad16f7e8649e94fb4fc21fe77e8310c060f61caaff8a","size":2},` + + `"layers":[{"mediaType":"application/vnd.oci.image.layer.v1.tar+gzip","digest":"` + digest + `","size":42}]}` + if string(one) != want { + t.Fatalf("the holder is\n%s\nwant\n%s", one, want) + } + if digestOf(emptyConfig) != emptyDigest { + t.Fatalf("the empty config's digest is %s, not %s", digestOf(emptyConfig), emptyDigest) + } +} + +func TestHoldBackfillsABareArchiveAndIsIdempotent(t *testing.T) { + // The archives published before this have no holder. Hold, run over every kept archive on + // every sweep, writes one the first time and nothing after. + m := &memRegistry{} + store := m.serve(t) + ctx := context.Background() + body := []byte("a theme") + reference := m.bare("shell/config", body) + + if held, err := store.Held(ctx, reference); err != nil || held { + t.Fatalf("a bare blob reads as held=%v (%v)", held, err) + } + wrote, err := store.Hold(ctx, reference) + if err != nil { + t.Fatal(err) + } + if !wrote { + t.Fatal("holding a bare archive wrote nothing") + } + _, holder := Holder(digestOf(body), int64(len(body))) + if _, ok := m.manifests["shell/config@"+holder]; !ok { + t.Fatalf("the store holds manifests %v; want %s", m.manifests, holder) + } + if held, err := store.Held(ctx, reference); err != nil || !held { + t.Fatalf("after holding, held=%v (%v)", held, err) + } + + writes := len(m.writes) + wrote, err = store.Hold(ctx, reference) + if err != nil { + t.Fatal(err) + } + if wrote || len(m.writes) != writes { + t.Fatalf("holding again wrote %v", m.writes[writes:]) + } +} + +func TestAnImageNeedsNoHolderAndAMissingArchiveIsGone(t *testing.T) { + m := &memRegistry{} + store := m.serve(t) + ctx := context.Background() + + // An image is its own manifest: nothing is asked. + wrote, err := store.Hold(ctx, catalogue.ArtifactStoreScheme+"web/app@sha256:"+strings.Repeat("b", 64)) + if err != nil || wrote || len(m.writes) != 0 { + t.Fatalf("holding an image wrote=%v err=%v writes=%v", wrote, err, m.writes) + } + // An archive the store does not have is a fact to report, not something to invent a holder for. + _, err = store.Hold(ctx, catalogue.ArtifactStoreScheme+"web/config/blobs/sha256:"+strings.Repeat("c", 64)) + if !errors.Is(err, Gone) { + t.Fatalf("holding a missing archive answered %v, want Gone", err) + } + // And a reference that is not the mesh's is refused as such. + if _, err := store.Hold(ctx, "docker.io/library/registry@sha256:abc"); !errors.Is(err, ErrNotOurs) { + t.Fatalf("holding a vendor's image answered %v, want ErrNotOurs", err) + } +} + +func TestLettingGoOfAnArchiveDeletesItsHolderFirst(t *testing.T) { + // A holder left behind would keep the bytes through every collection while the record said + // collected; the link deleted first and the holder failing after would be that exactly. + m := &memRegistry{} + store := m.serve(t) + ctx := context.Background() + body := []byte("an old theme") + reference := m.bare("shell/config", body) + if _, err := store.Hold(ctx, reference); err != nil { + t.Fatal(err) + } + m.writes = nil + + if err := store.LetGo(ctx, reference); err != nil { + t.Fatal(err) + } + _, holder := Holder(digestOf(body), int64(len(body))) + want := []string{ + "DELETE /v2/shell/config/manifests/" + holder, + "DELETE /v2/shell/config/blobs/" + digestOf(body), + } + if strings.Join(m.writes, "\n") != strings.Join(want, "\n") { + t.Fatalf("the store was asked\n%s\nwant\n%s", strings.Join(m.writes, "\n"), strings.Join(want, "\n")) + } + if len(m.manifests) != 0 { + t.Fatalf("a holder survived: %v", m.manifests) + } + // Asked again, the archive is already gone, which is the outcome wanted. + if err := store.LetGo(ctx, reference); !errors.Is(err, Gone) { + t.Fatalf("letting go twice answered %v, want Gone", err) + } +} + +func TestLettingGoOfAnUnheldArchiveStillDeletesIt(t *testing.T) { + // An archive published before holders and let go of before any sweep held it: the holder's + // delete answers 404, which is the outcome wanted, and the blob still goes. + m := &memRegistry{} + store := m.serve(t) + reference := m.bare("shell/config", []byte("never held")) + if err := store.LetGo(context.Background(), reference); err != nil { + t.Fatal(err) + } + if len(m.blobs) != 0 { + t.Fatalf("the blob survived: %v", m.blobs) + } +} diff --git a/internal/artifacts/live_registry_test.go b/internal/artifacts/live_registry_test.go new file mode 100644 index 0000000..c88b37f --- /dev/null +++ b/internal/artifacts/live_registry_test.go @@ -0,0 +1,130 @@ +package artifacts + +import ( + "bytes" + "context" + "fmt" + "io" + "net/http" + "os" + "os/exec" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// The registry's own collector keeps a held archive and takes a bare one (novox/hq issue 253, +// ADR 0189). +// +// Everything else here is asserted against a fake, which can only say what this side asks. This is +// the one question a fake cannot answer — what `registry garbage-collect` actually does with what +// this side wrote — and it is the whole of whether the store's nightly step may stop being a dry +// run. Against the very image the mesh's store runs: +// +// docker run -d --rm --name mesh-controller-registry -p 15000:5000 \ +// -e REGISTRY_STORAGE_DELETE_ENABLED=true registry:2.8.3 +// MESH_TEST_REGISTRY=127.0.0.1:15000 MESH_TEST_REGISTRY_CONTAINER=mesh-controller-registry \ +// go test -run Live ./internal/artifacts/ +// docker stop mesh-controller-registry +// +// Skipped without both variables: it needs a registry it may write to and collect, and a container +// to run the collector in. +func TestLiveTheRegistrysCollectorKeepsWhatIsHeldAndTakesWhatIsNot(t *testing.T) { + address := os.Getenv("MESH_TEST_REGISTRY") + container := os.Getenv("MESH_TEST_REGISTRY_CONTAINER") + if address == "" || container == "" { + t.Skip("no MESH_TEST_REGISTRY / MESH_TEST_REGISTRY_CONTAINER; see this test's comment for the registry to raise") + } + ctx := context.Background() + store := Store{Address: address} + run := time.Now().UnixNano() + + // Four archives in four repositories, each a different story. Distinct bytes per run, so a + // registry reused across runs cannot answer for an earlier one. + put := func(name string) (repository, digest string, body []byte) { + repository = fmt.Sprintf("live-%d/%s", run, name) + body = []byte(fmt.Sprintf("%s archive of run %d", name, run)) + digest = digestOf(body) + if err := store.putBlob(ctx, repository, digest, body); err != nil { + t.Fatal(err) + } + return repository, digest, body + } + reference := func(repository, digest string) string { + return catalogue.ArtifactStoreScheme + repository + "/blobs/" + digest + } + + // Published held — what PublishArchive now does. + heldRepo, heldDigest, heldBody := put("held") + if _, err := store.HoldBlob(ctx, heldRepo, heldDigest, int64(len(heldBody))); err != nil { + t.Fatalf("the registry refused a holder: %v", err) + } + // Published bare, as before, and never held: what the collector must take. + _, bareDigest, _ := put("bare") + // Published bare and then held by the sweep: the backfill. + backRepo, backDigest, backBody := put("backfilled") + if wrote, err := store.Hold(ctx, reference(backRepo, backDigest)); err != nil || !wrote { + t.Fatalf("backfilling wrote=%v: %v", wrote, err) + } + if held, err := store.Held(ctx, reference(backRepo, backDigest)); err != nil || !held { + t.Fatalf("after backfilling, held=%v: %v", held, err) + } + // Held, and then let go of by the sweep: holder first, then the link. + goneRepo, goneDigest, _ := put("let-go") + if _, err := store.Hold(ctx, reference(goneRepo, goneDigest)); err != nil { + t.Fatal(err) + } + if err := store.LetGo(ctx, reference(goneRepo, goneDigest)); err != nil { + t.Fatalf("letting go of a held archive: %v", err) + } + + collected, err := exec.CommandContext(ctx, "docker", "exec", container, + "registry", "garbage-collect", "/etc/docker/registry/config.yml").CombinedOutput() + if err != nil { + t.Fatalf("the collector failed: %v\n%s", err, collected) + } + t.Logf("the collector said:\n%s", lastLines(string(collected), 12)) + + // What is asserted is the bytes on the store's disk, not what the running server answers: the + // server caches blob descriptors in memory and can answer for a blob the collector removed. + onDisk := func(digest string) bool { + hex := strings.TrimPrefix(digest, "sha256:") + path := "/var/lib/registry/docker/registry/v2/blobs/sha256/" + hex[:2] + "/" + hex + "/data" + return exec.CommandContext(ctx, "docker", "exec", container, "test", "-f", path).Run() == nil + } + if !onDisk(heldDigest) { + t.Error("the collector took an archive published held") + } + if !onDisk(backDigest) { + t.Error("the collector took an archive the sweep backfilled a holder for") + } + if onDisk(bareDigest) { + t.Error("the collector kept a bare archive — then the holders prove nothing, and this test is wrong") + } + if onDisk(goneDigest) { + t.Error("the collector kept an archive the sweep let go of: its holder outlived its link") + } + // And what survived is still fetched exactly as machines fetch it: the blob, by digest. + for repository, want := range map[string][]byte{heldRepo: heldBody, backRepo: backBody} { + digest := digestOf(want) + response, err := http.Get(catalogue.Routed(reference(repository, digest), address)) + if err != nil { + t.Fatal(err) + } + got, _ := io.ReadAll(response.Body) + response.Body.Close() + if response.StatusCode != http.StatusOK || !bytes.Equal(got, want) { + t.Errorf("%s answered %s with %q after collection", repository, response.Status, got) + } + } +} + +func lastLines(s string, n int) string { + lines := strings.Split(strings.TrimSpace(s), "\n") + if len(lines) > n { + lines = lines[len(lines)-n:] + } + return strings.Join(lines, "\n") +} diff --git a/internal/artifacts/store.go b/internal/artifacts/store.go index ca80e65..688b933 100644 --- a/internal/artifacts/store.go +++ b/internal/artifacts/store.go @@ -1,7 +1,8 @@ // Package artifacts speaks to the mesh's artifact store over its own door. // // Only what the mesh needs that nothing else does: letting go of something it put there -// (novox/hq ADR 0189, issue 108). Pushing is the builder's, through the container runtime; reading +// (novox/hq ADR 0189, issue 108), and holding every archive it keeps by a manifest so the store's +// own collector does not take it (novox/hq issue 253). Pushing is the builder's, through the container runtime; reading // is every machine's, through its runtime. This is the one operation that belongs to the thing // holding the records, because it is the only one that is a decision rather than a transfer. package artifacts @@ -12,7 +13,6 @@ import ( "fmt" "net/http" "strings" - "time" "github.com/novox/mesh-controller/internal/catalogue" ) @@ -43,6 +43,9 @@ var ErrNotOurs = errors.New("not a reference into the mesh's artifact store") // an image, `…/blobs/sha256:…` for an archive — because that is the identity every record uses, // and composes the address here at the moment of use. // +// An archive is let go of in two deletes, its holder manifest and then the blob's link (novox/hq +// issue 253); an image in one. +// // Returns Gone when the store answers that it does not have it. That is not a failure: the sweep // wants the artifact absent, and it is. It is distinguished from success only so a caller can say // which of the two happened. @@ -66,17 +69,24 @@ func (s Store) LetGo(ctx context.Context, reference string) error { if err != nil { return err } - url := "http://" + s.Address + "/v2/" + repository + "/" + kind + "/" + digest + if kind == "blobs" { + // **An archive's holder goes before the archive** (novox/hq issue 253): a manifest left + // naming the blob would keep its bytes through every collection while the record said + // collected. Gone here means the store has no such blob, so there is nothing to let go. + if err := s.letGoOfHolder(ctx, repository, digest); err != nil { + return err + } + } + return s.remove(ctx, s.url(repository, kind, digest), reference) +} +// remove asks the store to delete what is at url. Gone when it has no such thing. +func (s Store) remove(ctx context.Context, url, what string) error { request, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil) if err != nil { return err } - client := s.HTTP - if client == nil { - client = &http.Client{Timeout: 30 * time.Second} - } - response, err := client.Do(request) + response, err := s.client().Do(request) if err != nil { return err } @@ -92,9 +102,9 @@ func (s Store) LetGo(ctx context.Context, reference string) error { return fmt.Errorf( "the artifact store refuses deletion: its server was started without it enabled "+ "(REGISTRY_STORAGE_DELETE_ENABLED), so nothing can be collected until the store "+ - "module is applied again (novox/hq ADR 0189). Asking about %s", reference) + "module is applied again (novox/hq ADR 0189). Asking about %s", what) default: - return fmt.Errorf("the artifact store answered %s for %s", response.Status, reference) + return fmt.Errorf("the artifact store answered %s for %s", response.Status, what) } } diff --git a/internal/artifacts/store_test.go b/internal/artifacts/store_test.go index ee09086..097b9a1 100644 --- a/internal/artifacts/store_test.go +++ b/internal/artifacts/store_test.go @@ -20,6 +20,17 @@ func fakeStore(t *testing.T, answer int) (Store, *[]string) { t.Helper() var asked []string server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/") { + // An archive's size, asked so its holder can be named (novox/hq issue 253). A store + // that does not have the thing does not have its blob either. + if answer == http.StatusNotFound { + w.WriteHeader(http.StatusNotFound) + return + } + w.Header().Set("Content-Length", "7") + w.WriteHeader(http.StatusOK) + return + } if r.Method != http.MethodDelete { t.Errorf("the store was asked %s %s; collecting is a delete", r.Method, r.URL.Path) } @@ -45,8 +56,14 @@ func TestAnImageAndAnArchiveAreAskedForAtTheirOwnEndpoints(t *testing.T) { if err := store.LetGo(ctx, archive); err != nil { t.Fatal(err) } - want := []string{"/v2/web/app/manifests/sha256:abc123", "/v2/web/config/blobs/sha256:def456"} - if len(*asked) != 2 || (*asked)[0] != want[0] || (*asked)[1] != want[1] { + // The archive's holder goes first, then the archive (novox/hq issue 253). + _, holder := Holder("sha256:def456", 7) + want := []string{ + "/v2/web/app/manifests/sha256:abc123", + "/v2/web/config/manifests/" + holder, + "/v2/web/config/blobs/sha256:def456", + } + if strings.Join(*asked, " ") != strings.Join(want, " ") { t.Fatalf("the store was asked %v; want %v", *asked, want) } } diff --git a/internal/builder/registry.go b/internal/builder/registry.go index 49ff966..219b4fe 100644 --- a/internal/builder/registry.go +++ b/internal/builder/registry.go @@ -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. diff --git a/internal/builder/registry_test.go b/internal/builder/registry_test.go index 20d4f00..1e3a440 100644 --- a/internal/builder/registry_test.go +++ b/internal/builder/registry_test.go @@ -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 // "@" + 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)) } } diff --git a/internal/inventory/collection.go b/internal/inventory/collection.go index 3f6cde4..684287a 100644 --- a/internal/inventory/collection.go +++ b/internal/inventory/collection.go @@ -3,6 +3,7 @@ package inventory import ( "context" "encoding/json" + "sort" "strings" "github.com/novox/mesh-controller/internal/catalogue" @@ -33,7 +34,8 @@ const KeptBuilds = 5 // - **a definition names it** — the reference appears in a module's recorded manifest, which is // what the mesh would hand a machine now. No age limit: this is the floor; // - **the mesh can still go back to it** — it is an artifact of one of the KeptBuilds most -// recent successful builds of its module; +// recent successful builds of a module the mesh still holds. A forgotten module keeps +// nothing beyond what a held definition names (novox/hq issue 253); // - it was already collected, in which case there is nothing left to do. // // Returned in a stated order so two runs over the same records ask for the same things in the @@ -109,12 +111,19 @@ func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error) return nil, err } - // The KeptBuilds most recent successful builds of each module, whole. + // The KeptBuilds most recent successful builds of each module the mesh still holds, whole. + // + // **Only a module the mesh still holds can be gone back to** (novox/hq issue 253). "Somewhere + // to return to" is a reason about a module's releases; a module that has been forgotten has + // no releases left to return between, and its build rows stay only as history. Without the + // join every module ever built kept five builds' artifacts for ever — and once the store's + // collector runs for real, what the keep set says is what the disk holds. recent, err := i.store.Pool().Query(ctx, `select made from ( - select made, row_number() over (partition by module order by at desc, id desc) as back - from build - where failed = '' and module is not null and module <> '' + select b.made, row_number() over (partition by b.module order by b.at desc, b.id desc) as back + from build b + join module m on m.name = b.module + where b.failed = '' and b.module is not null and b.module <> '' ) ranked where back <= $1`, KeptBuilds) if err != nil { return nil, err @@ -169,6 +178,28 @@ func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error) return keep, nil } +// KeptArchives is every archive the mesh keeps, in a stated order: the references the sweep must +// hold by a manifest before it lets anything go, and the ones an operator needs to read as all +// held before the store's collector is let loose (novox/hq issue 253, ADR 0189). +// +// Only references into the mesh's own store, and only blobs: an image is its own manifest, and a +// reference that is kept because nothing here can speak for it is not one the store can be asked +// about. +func (i *Inventory) KeptArchives(ctx context.Context) ([]string, error) { + keep, err := i.keptReferences(ctx) + if err != nil { + return nil, err + } + var out []string + for reference := range keep { + if path, ours := catalogue.InArtifactStore(reference); ours && strings.Contains(path, "/blobs/sha256:") { + out = append(out, reference) + } + } + sort.Strings(out) + return out, nil +} + // everyReferenceMade is every artifact reference any successful build recorded. func (i *Inventory) everyReferenceMade(ctx context.Context) ([]string, error) { rows, err := i.store.Pool().Query(ctx, diff --git a/internal/inventory/collection_test.go b/internal/inventory/collection_test.go index c7aaf77..e05a2a3 100644 --- a/internal/inventory/collection_test.go +++ b/internal/inventory/collection_test.go @@ -20,6 +20,19 @@ func ref(module, artifact string, n int) string { return fmt.Sprintf("%s%s/%s@sha256:%064x", catalogue.ArtifactStoreScheme, module, artifact, n) } +// holding registers a definition for each module that names no artifact, so the mesh holds the +// module and its recent builds are somewhere it can go back to — and nothing more. +func holding(t *testing.T, inv *Inventory, modules ...string) { + t.Helper() + for _, module := range modules { + m := catalogue.Manifest{Module: module, Version: "1"} + if err := inv.RegisterModule(context.Background(), m, + Source{Repository: "https://forge.invalid/" + module + ".git"}); err != nil { + t.Fatal(err) + } + } +} + // built records one successful build of a module publishing one image. func built(t *testing.T, inv *Inventory, id, module string, n int) string { t.Helper() @@ -35,6 +48,7 @@ func built(t *testing.T, inv *Inventory, id, module string, n int) string { func TestTheStoreKeepsTheRecentBuildsAndLetsGoOfTheRest(t *testing.T) { inv := fresh(t) ctx := context.Background() + holding(t, inv, "web") // Eight builds of one module, oldest first. Five are kept — the newest, and the four a // release that turns out wrong can be taken back to. @@ -98,6 +112,7 @@ func TestWhatHasBeenCollectedIsNotOfferedAgain(t *testing.T) { // time it runs, for ever — a number of requests that grows with the mesh's whole history. inv := fresh(t) ctx := context.Background() + holding(t, inv, "web") for i := 1; i <= 7; i++ { built(t, inv, fmt.Sprintf("b%02d", i), "web", i) } @@ -123,6 +138,7 @@ func TestWhatHasBeenCollectedIsNotOfferedAgain(t *testing.T) { func TestAFailedBuildNamesNothingToCollectAndEachModuleIsCountedOnItsOwn(t *testing.T) { inv := fresh(t) ctx := context.Background() + holding(t, inv, "web", "db") // A failed build published nothing, so it is neither kept nor collected — and it must not // count against the module's five. @@ -156,6 +172,7 @@ func TestAFailedBuildNamesNothingToCollectAndEachModuleIsCountedOnItsOwn(t *test func TestAnArtifactRecordedWithAnAddressIsOfferedAsTheMeshRecordsOne(t *testing.T) { inv := fresh(t) ctx := context.Background() + holding(t, inv, "tools") // The oldest build published the old way; five newer ones fill the module's five. old := aBuild("a00", "tools", "") @@ -190,3 +207,82 @@ func TestAnArtifactRecordedWithAnAddressIsOfferedAsTheMeshRecordsOne(t *testing. t.Fatalf("offered %v again after collecting it", again) } } + +// A forgotten module keeps nothing beyond what a held definition names (novox/hq issue 253). +// +// "Somewhere to go back to" is a reason about a module's releases, and a module the mesh no +// longer holds has none. Its build rows stay as history; its artifacts go — except one a module +// the mesh still holds names, which is the floor whatever built it. +func TestAForgottenModuleKeepsNothingAHeldDefinitionDoesNotName(t *testing.T) { + inv := fresh(t) + ctx := context.Background() + + // Three builds of a module that was never held, or was held and then forgotten: within its + // five, and kept for that reason until now. + var gone []string + for i := 1; i <= 3; i++ { + gone = append(gone, built(t, inv, fmt.Sprintf("o%02d", i), "old", 200+i)) + } + // A module the mesh holds, whose definition runs the forgotten module's newest image. + named := gone[2] + m := catalogue.Manifest{Module: "web", Version: "1", Resources: []map[string]any{{ + "id": "app", "type": "container", "name": "web", "image": named, + }}} + if err := inv.RegisterModule(ctx, m, Source{Repository: "https://forge.invalid/web.git"}); err != nil { + t.Fatal(err) + } + + go_, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + if len(go_) != 2 || go_[0] != gone[0] || go_[1] != gone[1] { + t.Fatalf("offered %v; want %v — a forgotten module's builds are no release to go back to, "+ + "and only what a held definition names stays", go_, gone[:2]) + } + + // And once the module is held again, its five are kept again. + holding(t, inv, "old") + again, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + if len(again) != 0 { + t.Fatalf("offered %v for a module the mesh holds, within its five", again) + } +} + +// The archives the mesh keeps are what the sweep holds before it lets anything go, and what an +// operator reads as all held before the store's collector is let loose (novox/hq issue 253). +func TestKeptArchivesAreTheKeptBlobsOnly(t *testing.T) { + inv := fresh(t) + ctx := context.Background() + holding(t, inv, "shell") + + archive := func(n int) string { + return fmt.Sprintf("%sshell/config/blobs/sha256:%064x", catalogue.ArtifactStoreScheme, n) + } + for i := 1; i <= 6; i++ { + b := aBuild(fmt.Sprintf("s%02d", i), "shell", "") + b.Made = []Artifact{ + {Name: "app", Kind: "image", Reference: ref("shell", "app", i)}, + {Name: "config", Kind: "archive", Reference: archive(i)}, + } + if err := inv.RecordBuild(ctx, b); err != nil { + t.Fatal(err) + } + } + kept, err := inv.KeptArchives(ctx) + if err != nil { + t.Fatal(err) + } + want := []string{archive(2), archive(3), archive(4), archive(5), archive(6)} + if len(kept) != len(want) { + t.Fatalf("kept archives %v; want the five recent ones and no images", kept) + } + for i := range want { + if kept[i] != want[i] { + t.Fatalf("kept archives %v; want %v", kept, want) + } + } +}