Files
mesh-controller/internal/artifacts/store.go
T
jochen 1cdd9df581
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
Let platform manifests go only on a confirmed collect, and copy again what a sweep took
An unrecorded index or a copy in progress can name a platform the records do not see, so
only a person's collect, after its dry run, takes an index's platforms, and only once every
kept index of each repository it touches was read. A copy missing a platform is copied
again, and a copy a build holds again is no longer recorded as collected (review of #144).
2026-10-08 15:32:10 +02:00

279 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, read whole before it starts: for every repository an
// image reference of `touching` is in, each kept manifest there and every platform manifest a kept
// index there names. Any read that fails is an error, and the sweep lets nothing go. A repository not
// read here is answered with an error, so an index in it is kept. A kept index the store no longer
// holds names nothing.
func (s Store) SpareKeptIndexes(ctx context.Context, kept, touching []string) (func(context.Context, string) (map[string]bool, error), error) {
imageIn := func(reference string) (repository, digest string, ok bool) {
path, ours := catalogue.InArtifactStore(reference)
if !ours {
return "", "", false
}
repository, kind, digest, err := split(path)
if err != nil || kind != "manifests" {
return "", "", false
}
return repository, digest, true
}
read := map[string]map[string]bool{}
for _, reference := range touching {
if repository, _, ok := imageIn(reference); ok {
read[repository] = map[string]bool{}
}
}
for _, reference := range kept {
repository, digest, ok := imageIn(reference)
spared, touched := read[repository]
if !ok || !touched {
continue
}
spared[digest] = true
children, err := s.Platforms(ctx, repository, digest)
if errors.Is(err, Gone) {
continue
}
if err != nil {
return nil, fmt.Errorf("reading the kept %s@%s: %w", repository, digest, err)
}
for _, c := range children {
spared[c] = true
}
}
return func(_ context.Context, repository string) (map[string]bool, error) {
spared, done := read[repository]
if !done {
return nil, fmt.Errorf("%s was not read before the sweep began", repository)
}
return spared, nil
}, 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...)
}