Files
mesh-controller/internal/artifacts/store.go
T
jochen 565974901a
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
Record the bases a build copies, keep them by the builds that stood on them, and copy each image once
A copied base was named only in what a build stood on, and nowhere when the build failed,
so the store's sweep could never let one go (hq issue 321). One repository per upstream
image stops each module asking the public registry for the same image again, and letting
an index go now takes its own platform manifests, which otherwise kept every byte. A
person can record the copies no record names through the new mirrors verb (hq ADR 0257).

The forge test fix is the same commit as on feat/plain-notifications: main fails without it.
2026-10-08 14:09:01 +02:00

269 lines
11 KiB
Go

// 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), 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
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"strings"
"github.com/novox/mesh-controller/internal/catalogue"
)
// Store is the artifact store at an address, as this machine reaches it.
type Store struct {
// Address is `host:port` — the store as the caller reaches it now, composed and never
// recorded (novox/hq 04-ISSUES/102).
Address string
// HTTP is the client used; nil is a client with a modest timeout.
HTTP *http.Client
// Spare says, for one repository, which manifests a kept index there names: the platform
// manifests an index being let go of must leave in place (novox/hq ADR 0257). Nil spares every
// platform manifest — an index is then let go of alone, as it was before.
Spare func(ctx context.Context, repository string) (map[string]bool, error)
}
// Gone is the answer when the store does not hold it: the outcome wanted, already true.
var Gone = errors.New("the store does not hold it")
// ErrNotOurs is a reference this sweep will not address: not the mesh's own, or naming nothing
// the store holds by digest.
//
// **A fact about the record, not about the store** (novox/hq issue 226). The two deserve opposite
// responses — skip one and go on, abandon the sweep for the other — and collapsing them into "an
// error" is how a cautious loop became one that did nothing while reporting the right number.
var ErrNotOurs = errors.New("not a reference into the mesh's artifact store")
// LetGo asks the store to drop one artifact the mesh recorded making.
//
// Takes a reference as the mesh records it — `artifact-store://<module>/<artifact>@sha256:…` for
// 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.
func (s Store) LetGo(ctx context.Context, reference string) error {
// **Strict, and deliberately** (novox/hq issue 226). Only a reference the mesh keeps in its
// own vocabulary is addressed here. `Recorded` would read `docker.io/library/registry@sha256:…`
// as the mesh's too — it cannot tell one registry host from another — so normalising belongs
// where the provenance is known, which is the sweep reading its own build records, not here
// where the only job is to refuse anything that is not plainly ours.
path, kept := catalogue.InArtifactStore(reference)
if !kept {
// Nothing the mesh put in its own store. Refused rather than attempted: composing a
// delete for a reference of unknown shape is how a sweep reaches something that is not
// the mesh's. Distinguished from a store that refuses, so a sweep skips this and goes on.
return fmt.Errorf("%w: %s", ErrNotOurs, reference)
}
if s.Address == "" {
return 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 err
}
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
}
}
if kind == "manifests" && s.Spare != nil {
// **An index's platform manifests go before the index** (novox/hq ADR 0257). Each was put
// under its own digest, so each is a manifest of the repository in its own right, and the
// store's collector keeps every manifest a repository holds: an index let go of alone frees
// no byte of the images it names. Before, not after, so a sweep cut short leaves the index
// to be asked about again, and with it the list of what is still to go.
if err := s.letGoOfPlatforms(ctx, repository, digest, reference); err != nil {
return err
}
}
return s.remove(ctx, s.url(repository, kind, digest), reference)
}
// letGoOfPlatforms deletes the manifests an index names, except those a kept index of the same
// repository also names. A manifest that is not an index names none.
func (s Store) letGoOfPlatforms(ctx context.Context, repository, digest, reference string) error {
children, err := s.Platforms(ctx, repository, digest)
if err != nil || len(children) == 0 {
// Gone: there is no index to read, and the delete that follows says so.
if errors.Is(err, Gone) {
return nil
}
return err
}
spared, err := s.Spare(ctx, repository)
if err != nil {
return fmt.Errorf("cannot say which platform manifests of %s a kept index names, so %s was kept: %w",
repository, reference, err)
}
for _, child := range children {
if spared[child] {
continue
}
if err := s.remove(ctx, s.url(repository, "manifests", child), reference+" (its "+child+")"); err != nil && !errors.Is(err, Gone) {
return err
}
}
return nil
}
// anyManifest is every manifest the store may hold, indexes included, so it answers with the
// document as it is.
var anyManifest = []string{
"application/vnd.oci.image.index.v1+json",
"application/vnd.docker.distribution.manifest.list.v2+json",
"application/vnd.oci.image.manifest.v1+json",
"application/vnd.docker.distribution.manifest.v2+json",
}
// Platforms is the manifests an index names, by digest; none for a manifest that is not an index.
// Gone when the store does not hold it.
func (s Store) Platforms(ctx context.Context, repository, digest string) ([]string, error) {
request, err := http.NewRequestWithContext(ctx, http.MethodGet, s.url(repository, "manifests", digest), nil)
if err != nil {
return nil, err
}
request.Header.Set("Accept", strings.Join(anyManifest, ", "))
response, err := s.client().Do(request)
if err != nil {
return nil, err
}
defer response.Body.Close()
switch response.StatusCode {
case http.StatusOK:
case http.StatusNotFound:
return nil, Gone
default:
return nil, fmt.Errorf("the artifact store answered %s when asked for %s@%s", response.Status, repository, digest)
}
var document struct {
Manifests []struct {
Digest string `json:"digest"`
} `json:"manifests"`
}
if err := json.NewDecoder(io.LimitReader(response.Body, 4<<20)).Decode(&document); err != nil {
return nil, fmt.Errorf("%s@%s is not a manifest: %w", repository, digest, err)
}
var out []string
for _, m := range document.Manifests {
if strings.HasPrefix(m.Digest, "sha256:") {
out = append(out, m.Digest)
}
}
return out, nil
}
// SpareKeptIndexes is a Spare for a sweep: the platform manifests the kept references of each
// repository name, read from the store once per repository and remembered for the sweep. A kept
// reference itself is spared too. A kept index the store does not hold names nothing.
func (s Store) SpareKeptIndexes(kept []string) func(ctx context.Context, repository string) (map[string]bool, error) {
byRepository := map[string][]string{}
for _, reference := range kept {
path, ours := catalogue.InArtifactStore(reference)
if !ours {
continue
}
repository, kind, digest, err := split(path)
if err != nil || kind != "manifests" {
continue
}
byRepository[repository] = append(byRepository[repository], digest)
}
read := map[string]map[string]bool{}
return func(ctx context.Context, repository string) (map[string]bool, error) {
if spared, done := read[repository]; done {
return spared, nil
}
spared := map[string]bool{}
for _, digest := range byRepository[repository] {
spared[digest] = true
children, err := s.Platforms(ctx, repository, digest)
if errors.Is(err, Gone) {
continue
}
if err != nil {
return nil, err
}
for _, c := range children {
spared[c] = true
}
}
read[repository] = spared
return spared, nil
}
}
// 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
}
response, err := s.client().Do(request)
if err != nil {
return err
}
defer response.Body.Close()
switch response.StatusCode {
case http.StatusAccepted, http.StatusOK, http.StatusNoContent:
return nil
case http.StatusNotFound:
return Gone
case http.StatusMethodNotAllowed:
// The registry was started without deletion enabled. Said plainly, because the remedy is
// a setting on the store's module and not anything about this artifact.
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", what)
default:
return fmt.Errorf("the artifact store answered %s for %s", response.Status, what)
}
}
// split reads a recorded path into the repository, which endpoint names the thing, and the digest.
//
// Two shapes, which are the two the mesh records: `<repository>@sha256:<hex>` is a manifest, and
// `<repository>/blobs/sha256:<hex>` is a blob.
func split(path string) (repository, kind, digest string, err error) {
if before, after, ok := strings.Cut(path, "@sha256:"); ok {
return before, "manifests", "sha256:" + after, nil
}
if before, after, ok := strings.Cut(path, "/blobs/sha256:"); ok {
return before, "blobs", "sha256:" + after, nil
}
return "", "", "", fmt.Errorf("%w: %q names nothing the store holds by digest", ErrNotOurs, path)
}
// HoldsManifest is whether the store holds the manifest a recorded image reference names.
func (s Store) HoldsManifest(ctx context.Context, reference string) (bool, error) {
path, ours := catalogue.InArtifactStore(reference)
if !ours {
return false, fmt.Errorf("%w: %s", ErrNotOurs, reference)
}
repository, kind, digest, err := split(path)
if err != nil {
return false, err
}
if kind != "manifests" {
return false, fmt.Errorf("%w: %s names a blob, not a manifest", ErrNotOurs, reference)
}
return s.has(ctx, s.url(repository, kind, digest), anyManifest...)
}