Merge pull request 'Record the bases a build copies, keep them by the builds that stood on them, copy each image once (hq ADR 0257, issue 321)' (#144) from feat/a-mirror-is-recorded into main

This commit was merged in pull request #144.
This commit is contained in:
2026-10-08 14:02:47 +00:00
27 changed files with 1968 additions and 36 deletions
+2
View File
@@ -264,6 +264,8 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri
request.Repository, request.Path, request.Ref, workspace, request.Held, npmrc, request.Repository, request.Path, request.Ref, workspace, request.Held, npmrc,
forgeFrom(), say, request.Seats) forgeFrom(), say, request.Seats)
} }
// What it copied into the store, whatever became of the build (novox/hq ADR 0257).
result.Mirrored = built.Mirrored
// Only a build the kill ended: the kill came before its work did. One that finished — built, or // Only a build the kill ended: the kill came before its work did. One that finished — built, or
// failed on its own — in the moment the kill arrived says what it did, and the kill is refused. // failed on its own — in the moment the kill arrived says what it did, and the kill is refused.
killed := false killed := false
+3
View File
@@ -160,6 +160,9 @@ func buildFrom(result link.BuildResult) inventory.Build {
for _, ref := range result.Against { for _, ref := range result.Against {
kept.Against = append(kept.Against, catalogue.Recorded(ref)) kept.Against = append(kept.Against, catalogue.Recorded(ref))
} }
for _, ref := range result.Mirrored {
kept.Mirrored = append(kept.Mirrored, catalogue.Recorded(ref))
}
for _, r := range result.Read { for _, r := range result.Read {
kept.Read = append(kept.Read, inventory.ReadRepository{Repository: r.Repository, Ref: r.Ref}) kept.Read = append(kept.Read, inventory.ReadRepository{Repository: r.Repository, Ref: r.Ref})
} }
+50
View File
@@ -31,6 +31,11 @@ type sweepBounds struct {
most int most int
// budget is the longest it keeps its caller waiting. // budget is the longest it keeps its caller waiting.
budget time.Duration budget time.Duration
// platforms is whether an index is let go of with its own platform manifests (novox/hq ADR 0257).
// Only a sweep a person confirmed, after its dry run listed what stays and names a platform: the
// records cannot see an unrecorded index, or a copy in progress, naming the same platform, and the
// store's tools can.
platforms bool
} }
// afterBuild are the bounds of the sweep inside somebody's build. // afterBuild are the bounds of the sweep inside somebody's build.
@@ -54,6 +59,9 @@ type sweepResult struct {
LetGo []string LetGo []string
// Skipped is how many references the sweep will not address (novox/hq issue 226). // Skipped is how many references the sweep will not address (novox/hq issue 226).
Skipped int Skipped int
// Indexes is how many eligible indexes a sweep a person did not confirm left for one that is: an
// index goes only with its platform manifests (novox/hq ADR 0257 §4).
Indexes int
// Left is how many eligible references were not asked about this time. // Left is how many eligible references were not asked about this time.
Left int Left int
// Bounded is whether it stopped at its bounds rather than at a refusal. // Bounded is whether it stopped at its bounds rather than at a refusal.
@@ -76,6 +84,29 @@ func sweep(ctx context.Context, inv *inventory.Inventory, store artifacts.Store,
// as builds come, and is two HEADs each once done. A kept archive that could not be held stops // 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 // 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. // collector's safety rests on, and a store refusing a hold would refuse the deletes too.
// **An index goes with its platform manifests only when a person confirmed it** (novox/hq ADR
// 0257), and a kept index keeps its own: which platform manifests the kept indexes name is read
// from the store for every repository the sweep will touch, before the first delete. A single
// read that fails stops the sweep before it deletes anything, as a kept archive that cannot be
// held does. The sweep after a build leaves an eligible index for such a collect, below.
store.Spare = nil
if bounds.platforms {
if inv == nil {
r.Stopped = "no records to read which platform manifests a kept index names, so nothing was let go"
r.Left = len(references)
return r
}
images, err := inv.KeptImages(ctx)
if err == nil {
store.Spare, err = store.SpareKeptIndexes(within, images, references)
}
if err != nil {
r.Stopped = fmt.Sprintf("which platform manifests a kept index names could not be read, so nothing was let go: %v", err)
r.Left = len(references)
return r
}
}
wrote, missing, err := holdKept(within, store, kept) wrote, missing, err := holdKept(within, store, kept)
r.HoldersWritten, r.Missing = wrote, missing r.HoldersWritten, r.Missing = wrote, missing
if err != nil { if err != nil {
@@ -96,6 +127,22 @@ func sweep(ctx context.Context, inv *inventory.Inventory, store artifacts.Store,
} }
break break
} }
if !bounds.platforms {
// **An index waits for a confirmed collect** (novox/hq ADR 0257 §4). Let go of alone,
// its platform manifests would stay in the store and the record would say collected, so
// no later sweep would offer it again and they would stay for ever. Left eligible, the
// next collect a person confirms takes it with them.
index, err := store.IsIndex(within, reference)
if err != nil && !errors.Is(err, artifacts.Gone) && !errors.Is(err, artifacts.ErrNotOurs) {
r.Stopped = fmt.Sprintf("the artifact store could not say whether %s is an index, so nothing more was asked of it: %v", reference, err)
r.Left = len(references) - i
break
}
if index {
r.Indexes++
continue
}
}
err := store.LetGo(within, reference) err := store.LetGo(within, reference)
if err == nil || errors.Is(err, artifacts.Gone) { if err == nil || errors.Is(err, artifacts.Gone) {
// Gone is the outcome wanted, already true. Recorded so the next sweep does not ask // Gone is the outcome wanted, already true. Recorded so the next sweep does not ask
@@ -187,6 +234,9 @@ func collect(ctx context.Context, inv *inventory.Inventory) {
if r.Skipped > 0 { if r.Skipped > 0 {
fmt.Fprintf(os.Stderr, "%d artifact(s) the sweep will not address were skipped\n", r.Skipped) fmt.Fprintf(os.Stderr, "%d artifact(s) the sweep will not address were skipped\n", r.Skipped)
} }
if r.Indexes > 0 {
fmt.Fprintf(os.Stderr, "%d eligible index(es) wait for a collect a person confirms, which takes their platforms too\n", r.Indexes)
}
} }
// holdKept holds every kept archive by its manifest, stopping at the first refusal by the store. // holdKept holds every kept archive by its manifest, stopping at the first refusal by the store.
+4
View File
@@ -86,6 +86,10 @@ var handActVerbs = []handActVerb{
// a moment the person chose (ADR 0251) — never a repair. // a moment the person chose (ADR 0251) — never a repair.
{Verb: "collect", Decision: "letting the store go of what the records keep for no reason, now rather " + {Verb: "collect", Decision: "letting the store go of what the records keep for no reason, now rather " +
"than at the next build, is a person's word (ADR 0251)"}, "than at the next build, is a person's word (ADR 0251)"},
// Recording a copy no record names as the mesh's: a person's reading of the store, never a repair
// (ADR 0257).
{Verb: "mirrors", Decision: "which copies in the store no record names are the mesh's is a person's " +
"word (ADR 0257)"},
{Verb: "bus upgrade", Decision: "the bus is never rolled by the mesh: replacing it is a planned step a " + {Verb: "bus upgrade", Decision: "the bus is never rolled by the mesh: replacing it is a planned step a " +
"person starts (ADR 0236)"}, "person starts (ADR 0236)"},
{Verb: "upgrade release-backlog", Decision: "after a release plan failed, the next opens only when a " + {Verb: "upgrade release-backlog", Decision: "after a release plan failed, the next opens only when a " +
+2
View File
@@ -100,6 +100,8 @@ func run() error {
return collectCommand(ctx, args[1:]) return collectCommand(ctx, args[1:])
case "images": case "images":
return imagesCommand(ctx, args[1:]) return imagesCommand(ctx, args[1:])
case "mirrors":
return mirrorsCommand(ctx, args[1:])
case "plans": case "plans":
return plansCommand(ctx, args[1:]) return plansCommand(ctx, args[1:])
case "delivery": case "delivery":
+237
View File
@@ -0,0 +1,237 @@
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"os"
"regexp"
"strings"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// The bases the mesh copied into its artifact store, and the copies no record names (novox/hq ADR 0257).
//
// A build that stands on an image published elsewhere copies it in first (ADR 0096, 0097). Since ADR 0257
// each copy is recorded where it is made, and a copy is kept while a kept build stood on it. Copies made
// before then by a build that failed — or before the records kept what a build stood on — are in the store
// with no record naming them, and ADR 0189 §3 keeps the sweep away from anything no record names. This
// verb is how a person says such a copy is the mesh's: it records it, and the sweep then decides as it
// does for every other copy. It deletes nothing.
// mirrorRepository is a repository only the mirror writes: one image's repository, or a module's own
// repository of a base from before ADR 0257.
var mirrorRepository = regexp.MustCompile(`^(upstream/[a-z0-9][a-z0-9._/-]*|[a-z0-9][a-z0-9._-]*/on-[a-z0-9_]+)$`)
// mirrorReference reads a reference a person gave into its recorded form, refusing anything that is not
// a manifest in a repository the mirror writes.
func mirrorReference(given string) (string, error) {
reference := catalogue.Recorded(strings.TrimSpace(given))
path, ours := catalogue.InArtifactStore(reference)
if !ours {
return "", fmt.Errorf("%q is not a reference into the artifact store (artifact-store://<repository>@sha256:<hex>)", given)
}
repository, digest, ok := strings.Cut(path, "@sha256:")
if !ok || len(digest) != 64 || strings.Trim(digest, "0123456789abcdef") != "" {
return "", fmt.Errorf("%q names no manifest by digest", given)
}
if !mirrorRepository.MatchString(repository) {
return "", fmt.Errorf("%q is in %s, which is not a repository the mirror writes "+
"(upstream/<host>/<path>, or <module>/on-<argument>)", given, repository)
}
return reference, nil
}
type mirrorEntry struct {
Reference string `json:"reference"`
State string `json:"state"`
Why []string `json:"why,omitempty"`
}
type mirrorsAnswer struct {
DryRun bool `json:"dry_run,omitempty"`
// Mirrors is every copy the records hold that is not yet let go of.
Mirrors []mirrorEntry `json:"mirrors"`
Counts struct {
Kept int `json:"kept"`
Eligible int `json:"eligible"`
Collected int `json:"collected"`
} `json:"counts"`
// Recorded, AlreadyRecorded and Refused answer a record: what was (or, a dry run, would be)
// recorded, what the records already held, and what was refused, with why.
Recorded []string `json:"recorded,omitempty"`
// Then says what recording does: a recorded copy is eligible, and the next sweep lets it go.
Then string `json:"then,omitempty"`
AlreadyRecorded []string `json:"already_recorded,omitempty"`
Refused map[string]string `json:"refused,omitempty"`
}
// mirrorsOf is the copies among the recorded states.
func mirrorsOf(states []inventory.ArtifactState, mirrored map[string]bool) mirrorsAnswer {
a := mirrorsAnswer{Mirrors: []mirrorEntry{}}
for _, s := range states {
if !mirrored[s.Reference] {
continue
}
switch s.State {
case inventory.ArtifactKept:
a.Counts.Kept++
case inventory.ArtifactEligible:
a.Counts.Eligible++
case inventory.ArtifactCollected:
a.Counts.Collected++
continue
}
a.Mirrors = append(a.Mirrors, mirrorEntry{Reference: s.Reference, State: s.State, Why: s.Why})
}
return a
}
// splitReferences reads references separated by spaces or commas.
func splitReferences(given string) []string {
return strings.FieldsFunc(given, func(r rune) bool { return r == ',' || r == ' ' || r == '\n' || r == '\t' })
}
// recordMirrors decides, for each reference given, whether it may be recorded: in a repository the
// mirror writes, not already recorded, and held by the store. A real run records those.
func recordMirrors(ctx context.Context, inv *inventory.Inventory, store artifacts.Store, given []string,
known map[string]bool, why string, real bool) (mirrorsAnswer, error) {
a := mirrorsAnswer{DryRun: !real, Mirrors: []mirrorEntry{}, Refused: map[string]string{}}
var record []string
seen := map[string]bool{}
for _, g := range given {
reference, err := mirrorReference(g)
if err != nil {
a.Refused[g] = err.Error()
continue
}
if seen[reference] {
continue
}
seen[reference] = true
if known[reference] {
a.AlreadyRecorded = append(a.AlreadyRecorded, reference)
continue
}
if store.Address == "" {
a.Refused[g] = "this mesh has no artifact store on its network to ask whether it holds this"
continue
}
held, err := store.HoldsManifest(ctx, reference)
switch {
case err != nil:
a.Refused[g] = err.Error()
continue
case !held:
a.Refused[g] = "the artifact store does not hold it"
continue
}
record = append(record, reference)
}
a.Recorded = record
if len(record) > 0 {
verb := "would become"
if real {
verb = "became"
}
a.Then = fmt.Sprintf("%d cop%s %s eligible: the next sweep (after any build, or collect) lets each go "+
"unless one of the five kept builds of a module the mesh holds stood on it", len(record),
map[bool]string{true: "y", false: "ies"}[len(record) == 1], verb)
}
if real && len(record) > 0 {
if err := inv.RecordMirrored(ctx, "", why, record); err != nil {
return a, fmt.Errorf("recording %d copies: %w", len(record), err)
}
}
return a, nil
}
func mirrorsCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("mirrors", flag.ContinueOnError)
asJSON := set.Bool("json", false, "answer as JSON")
record := set.String("record", "", "references of copies no record names, separated by spaces or commas, to record as the mesh's")
confirm := set.Bool("confirm", false, "record them, rather than only say what would be recorded")
f := addHandActFlags(set)
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 0 {
return errors.New("mirrors [--json] [--record <references> [--confirm --why <text>]]")
}
if *confirm {
if strings.TrimSpace(*record) == "" {
return errors.New("mirrors: --confirm records what --record names, and it names nothing. Nothing was done")
}
if err := f.require("mirrors"); err != nil {
return err
}
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
states, err := inv.Artifacts(ctx)
if err != nil {
return err
}
known, err := inv.Mirrored(ctx)
if err != nil {
return err
}
var answer mirrorsAnswer
if strings.TrimSpace(*record) == "" {
answer = mirrorsOf(states, known)
} else {
shelf, err := inv.Catalogue(ctx)
if err != nil {
return err
}
address, err := artifactStoreAddress(ctx, inv, shelf, "")
if err != nil {
return err
}
answer, err = recordMirrors(ctx, inv, artifacts.Store{Address: address}, splitReferences(*record), known,
strings.TrimSpace(*f.why), *confirm)
if err != nil {
return err
}
// The hand act is logged once the records say what it did, never before: an act that
// failed half-way is not logged as done.
if *confirm && len(answer.Recorded) > 0 {
f.record(ctx, "mirrors", append([]string{"--record"}, answer.Recorded...))
}
}
if *asJSON {
encoder := json.NewEncoder(os.Stdout)
encoder.SetIndent("", " ")
return encoder.Encode(answer)
}
if answer.DryRun && len(answer.Recorded) > 0 {
fmt.Println("a dry run: nothing was recorded (--confirm --why <text> to record)")
}
if answer.Then != "" {
fmt.Println(answer.Then)
}
for _, r := range answer.Recorded {
fmt.Printf(" record %s\n", r)
}
for _, r := range answer.AlreadyRecorded {
fmt.Printf(" already %s\n", r)
}
for g, why := range answer.Refused {
fmt.Printf(" refused %s: %s\n", g, why)
}
for _, m := range answer.Mirrors {
fmt.Printf(" %-9s %s %s\n", m.State, m.Reference, strings.Join(m.Why, ", "))
}
return nil
}
+298
View File
@@ -0,0 +1,298 @@
package main
import (
"fmt"
"net/http"
"net/http/httptest"
"slices"
"strings"
"sync"
"testing"
"time"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// The bases the mesh copied in, and recording a copy no record names (novox/hq ADR 0257).
func TestTheMirrorsVerbComposesItsCommandAndRefusesWhatItDoesNotTake(t *testing.T) {
r := ref("route-proxy", "on-go_base", 1)
for _, c := range []struct {
args map[string]any
want string
}{
{map[string]any{}, "mirrors --json"},
{map[string]any{"record": r}, "mirrors --json --record " + r},
{map[string]any{"record": r, "confirm": "true", "why": "copied before copies were recorded"},
"mirrors --json --record " + r + " --confirm --why copied before copies were recorded"},
} {
argv, err := argvFor("mirrors", c.args)
if err != nil || strings.Join(argv, " ") != c.want {
t.Errorf("mirrors %v: %q %v, want %q", c.args, argv, err, c.want)
}
}
for _, args := range []map[string]any{
{"record": r, "confirm": "true"}, // recorded without a why
{"record": r, "why": "x"}, // a why for a dry run
{"confirm": "true", "why": "x"}, // nothing to record
{"record": r, "confirm": "yes", "why": "x"}, // a switch is true or false
{"delete": "true"},
} {
if argv, err := argvFor("mirrors", args); err == nil {
t.Errorf("mirrors %v was composed as %q", args, argv)
}
}
if repairingCommand([]string{"mirrors", "--json", "--record", r, "--confirm", "--why", "x"}) != "mirrors" {
t.Error("recording a copy through the generic verb would go unrecorded")
}
if repairingCommand([]string{"mirrors", "--json", "--record", r}) != "" {
t.Error("a dry run was taken for an act by hand")
}
if !personsDecision(link.HandAct{Verb: "mirrors"}) {
t.Error("recording a copy would count toward a healer the mesh lacks")
}
}
func TestOnlyACopyInARepositoryTheMirrorWritesIsTaken(t *testing.T) {
for given, ok := range map[string]bool{
ref("route-proxy", "on-go_base", 1): true,
catalogue.ArtifactStoreScheme + "upstream/docker.io/library/golang@sha256:" + strings.Repeat("a", 64): true,
ref("route-proxy", "server", 1): false, // a module's own image
catalogue.ArtifactStoreScheme + "novox/invoicing-api@sha256:" + strings.Repeat("a", 64): false,
archiveRef("web", "on-tools", 1): false, // a blob
catalogue.ArtifactStoreScheme + "web/on-x@sha256:abc": false, // not a whole digest
"docker.io/library/golang@sha256:" + strings.Repeat("a", 64): false,
} {
if _, err := mirrorReference(given); (err == nil) != ok {
t.Errorf("%s: taken %v, want %v (%v)", given, err == nil, ok, err)
}
}
}
// heldStore answers a manifest HEAD with 200 for the digests it holds, and records every request.
func heldStore(t *testing.T, holds ...string) (artifacts.Store, *[]string) {
t.Helper()
var asked []string
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
asked = append(asked, r.Method+" "+r.URL.Path)
if r.Method == http.MethodHead && slices.ContainsFunc(holds, func(d string) bool { return strings.HasSuffix(r.URL.Path, d) }) {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
}))
t.Cleanup(server.Close)
return artifacts.Store{Address: strings.TrimPrefix(server.URL, "http://")}, &asked
}
func TestRecordingACopyTakesWhatTheStoreHoldsAndDeletesNothing(t *testing.T) {
inv := inventory.ForTest(t)
held := ref("route-proxy", "on-go_base", 1)
absent := ref("route-proxy", "on-go_base", 2)
known := ref("nats", "on-nats_base", 3)
notACopy := ref("route-proxy", "server", 4)
if err := inv.RecordMirrored(t.Context(), "b1", "", []string{known}); err != nil {
t.Fatal(err)
}
store, asked := heldStore(t, fmt.Sprintf("%064x", 1))
given := []string{held, absent, known, notACopy}
mirrored, err := inv.Mirrored(t.Context())
if err != nil {
t.Fatal(err)
}
// A dry run says, and records nothing.
dry, err := recordMirrors(t.Context(), inv, store, given, mirrored, "", false)
if err != nil {
t.Fatal(err)
}
if !dry.DryRun || !slices.Equal(dry.Recorded, []string{held}) || !slices.Equal(dry.AlreadyRecorded, []string{known}) ||
len(dry.Refused) != 2 || dry.Refused[absent] == "" || dry.Refused[notACopy] == "" {
t.Fatalf("a dry run answered %+v", dry)
}
if !strings.Contains(dry.Then, "would become eligible") || !strings.Contains(dry.Then, "next sweep") {
t.Fatalf("a dry run did not say what recording does: %q", dry.Then)
}
if again, _ := inv.Mirrored(t.Context()); again[held] {
t.Fatal("a dry run recorded a copy")
}
// A real one records what the store holds, with the person's why, and the sweep may now decide.
real, err := recordMirrors(t.Context(), inv, store, given, mirrored, "copied before copies were recorded", true)
if err != nil {
t.Fatal(err)
}
if !slices.Equal(real.Recorded, []string{held}) {
t.Fatalf("recorded %v", real.Recorded)
}
left, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Contains(left, held) {
t.Fatalf("a recorded copy no kept build stood on is not offered to collect: %v", left)
}
for _, r := range *asked {
if !strings.HasPrefix(r, "HEAD ") {
t.Fatalf("recording a copy asked the store %s", r)
}
}
}
// indexRegistry holds indexes and platforms of one image's repository, answers GETs with their
// documents, refuses GETs for the digests in fail, and records every delete.
type indexRegistry struct {
mu sync.Mutex
indexes map[string][]string
held map[string]bool
fail map[string]bool
deleted []string
}
func (f *indexRegistry) serve(t *testing.T) artifacts.Store {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
f.mu.Lock()
defer f.mu.Unlock()
digest := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]
switch r.Method {
case http.MethodGet:
if f.fail[digest] {
w.WriteHeader(http.StatusInternalServerError)
return
}
if children, ok := f.indexes[digest]; ok && f.held[digest] {
var named []string
for _, c := range children {
named = append(named, `{"mediaType":"application/vnd.oci.image.manifest.v1+json","digest":"`+c+`"}`)
}
_, _ = w.Write([]byte(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.index.v1+json","manifests":[` +
strings.Join(named, ",") + `]}`))
return
}
if f.held[digest] {
_, _ = w.Write([]byte(`{"schemaVersion":2,"layers":[]}`))
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodHead:
if f.held[digest] {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodDelete:
if !f.held[digest] {
w.WriteHeader(http.StatusNotFound)
return
}
delete(f.held, digest)
f.deleted = append(f.deleted, digest)
w.WriteHeader(http.StatusAccepted)
default:
w.WriteHeader(http.StatusBadRequest)
}
}))
t.Cleanup(server.Close)
return artifacts.Store{Address: strings.TrimPrefix(server.URL, "http://")}
}
func copyDigest(n int) string { return fmt.Sprintf("sha256:%064x", n) }
// twoCopies records, in one image's repository, a kept copy (stood on by a held module's build) naming
// platforms 11 and 12, and an eligible one (a failed build's) naming 12 and 13.
func twoCopies(t *testing.T, inv *inventory.Inventory) (kept, eligible string, reg *indexRegistry) {
t.Helper()
const repository = "upstream/docker.io/library/golang"
kept = catalogue.ArtifactStoreScheme + repository + "@" + copyDigest(1)
eligible = catalogue.ArtifactStoreScheme + repository + "@" + copyDigest(2)
if err := inv.RegisterModule(t.Context(), catalogue.Manifest{Module: "proxy", Version: "1"},
inventory.Source{Repository: "https://forge.invalid/proxy.git"}); err != nil {
t.Fatal(err)
}
failed := inventory.Build{ID: "f1", Repository: "https://forge.invalid/proxy.git", On: "a-build-machine",
Failed: "a recipe refused", Mirrored: []string{eligible}}
if err := inv.RecordBuild(t.Context(), failed); err != nil {
t.Fatal(err)
}
ok := inventory.Build{ID: "b1", Repository: "https://forge.invalid/proxy.git", Module: "proxy", On: "a-build-machine",
Commit: "c0ffee", Made: []inventory.Artifact{{Name: "app", Kind: "image", Reference: ref("proxy", "app", 7)}},
Against: []string{kept}, Mirrored: []string{kept}}
if err := inv.RecordBuild(t.Context(), ok); err != nil {
t.Fatal(err)
}
reg = &indexRegistry{
indexes: map[string][]string{copyDigest(1): {copyDigest(11), copyDigest(12)}, copyDigest(2): {copyDigest(12), copyDigest(13)}},
held: map[string]bool{copyDigest(1): true, copyDigest(2): true, copyDigest(11): true, copyDigest(12): true, copyDigest(13): true},
fail: map[string]bool{},
}
return kept, eligible, reg
}
// Through the records: a confirmed collect lets the eligible index go with the platform only it names,
// and leaves the platform the kept index names.
func TestAConfirmedCollectLetsAnIndexGoWithItsOwnPlatformsOnly(t *testing.T) {
inv := inventory.ForTest(t)
_, eligible, reg := twoCopies(t, inv)
store := reg.serve(t)
references, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Equal(references, []string{eligible}) {
t.Fatalf("offered %v, want the failed build's copy", references)
}
a := runCollect(t.Context(), inv, store, references, nil, true,
sweepBounds{most: 10, budget: 5 * time.Second, platforms: true})
if !slices.Equal(a.LetGo, []string{eligible}) || a.Stopped != "" {
t.Fatalf("let go %v, stopped %q", a.LetGo, a.Stopped)
}
if !slices.Equal(reg.deleted, []string{copyDigest(13), copyDigest(2)}) {
t.Fatalf("deleted %v; want its own platform, then the index", reg.deleted)
}
}
// The sweep after a build leaves an eligible index in place and eligible: let go of alone, its
// platforms would stay for ever under a record that says collected. A later collect a person confirms
// takes it with the platform only it names. An image the same sweep reaches is let go of as before.
func TestTheSweepAfterABuildLeavesAnIndexForAConfirmedCollect(t *testing.T) {
inv := inventory.ForTest(t)
_, eligible, reg := twoCopies(t, inv)
reg.held[copyDigest(5)] = true
image := catalogue.ArtifactStoreScheme + "upstream/docker.io/library/golang@" + copyDigest(5)
store := reg.serve(t)
r := sweep(t.Context(), inv, store, []string{eligible, image}, nil, afterBuild)
if r.Indexes != 1 || !slices.Equal(r.LetGo, []string{image}) || !slices.Equal(reg.deleted, []string{copyDigest(5)}) {
t.Fatalf("after a build: %d indexes left, let go %v, deleted %v; want the index left and the image gone",
r.Indexes, r.LetGo, reg.deleted)
}
left, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Equal(left, []string{eligible}) {
t.Fatalf("after the sweep the records offer %v; want the index still eligible", left)
}
a := runCollect(t.Context(), inv, store, left, nil, true,
sweepBounds{most: 10, budget: 5 * time.Second, platforms: true})
if !slices.Equal(a.LetGo, []string{eligible}) ||
!slices.Equal(reg.deleted, []string{copyDigest(5), copyDigest(13), copyDigest(2)}) {
t.Fatalf("a confirmed collect let go %v, deleted %v; want the index with its own platform", a.LetGo, reg.deleted)
}
}
// A kept index the store will not answer for stops a confirmed collect before its first delete.
func TestASpareListThatCannotBeReadStopsTheSweepBeforeAnyDelete(t *testing.T) {
inv := inventory.ForTest(t)
_, eligible, reg := twoCopies(t, inv)
reg.fail[copyDigest(1)] = true
store := reg.serve(t)
a := runCollect(t.Context(), inv, store, []string{eligible}, nil, true,
sweepBounds{most: 10, budget: 5 * time.Second, platforms: true})
if len(a.LetGo) != 0 || len(reg.deleted) != 0 || !strings.Contains(a.Stopped, "nothing was let go") {
t.Fatalf("let go %v, deleted %v, stopped %q", a.LetGo, reg.deleted, a.Stopped)
}
}
+1 -1
View File
@@ -209,7 +209,7 @@ func collectCommand(ctx context.Context, args []string) error {
if *most < 1 { if *most < 1 {
return fmt.Errorf("collect: --most is at least 1, not %d", *most) return fmt.Errorf("collect: --most is at least 1, not %d", *most)
} }
bounds := sweepBounds{most: min(*most, collectMostAtMost), budget: collectBudget} bounds := sweepBounds{most: min(*most, collectMostAtMost), budget: collectBudget, platforms: true}
open, err := openStores(ctx) open, err := openStores(ctx)
if err != nil { if err != nil {
@@ -133,6 +133,10 @@ func (f *fakeStore) serve(t *testing.T) artifacts.Store {
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
case r.Method == http.MethodHead: case r.Method == http.MethodHead:
w.WriteHeader(http.StatusNotFound) w.WriteHeader(http.StatusNotFound)
case r.Method == http.MethodGet && strings.Contains(r.URL.Path, "/manifests/"):
// An image, not an index: it names no platform manifests.
w.Header().Set("Content-Type", "application/vnd.oci.image.manifest.v1+json")
_, _ = w.Write([]byte(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","layers":[]}`))
case r.Method == http.MethodPost: case r.Method == http.MethodPost:
w.Header().Set("Location", "/upload/x?state=1") w.Header().Set("Location", "/upload/x?state=1")
w.WriteHeader(http.StatusAccepted) w.WriteHeader(http.StatusAccepted)
+27 -1
View File
@@ -605,6 +605,30 @@ func (a *verbArguments) commandLine() ([]string, error) {
return nil, err return nil, err
} }
return []string{"images", str("node"), "--json"}, nil return []string{"images", str("node"), "--json"}, nil
case "mirrors":
// novox/hq ADR 0257.
argv := []string{"mirrors", "--json"}
record := str("record")
why := str("why")
if record == "" {
if on("confirm") || why != "" {
return nil, errors.New("mirrors takes confirm and why only with record: there is nothing " +
"else it records. Nothing was done")
}
return argv, nil
}
argv = append(argv, "--record", record)
if !on("confirm") {
if why != "" {
return nil, errors.New("mirrors takes why only with confirm: without confirm it is a dry run, " +
"and a reason for nothing would be recorded nowhere. Nothing was done")
}
return argv, nil
}
if err := need("why"); err != nil {
return nil, fmt.Errorf("%w: recording a copy as the mesh's is a hand act, which says why. Nothing was done", err)
}
return append(argv, "--confirm", "--why", why), nil
case "data": case "data":
argv := []string{"data", "--json"} argv := []string{"data", "--json"}
if m := str("machine"); m != "" { if m := str("machine"); m != "" {
@@ -804,7 +828,7 @@ func (a *verbArguments) commandLine() ([]string, error) {
var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true, "collection": true, var jsonVerbs = map[string]bool{"status": true, "seats": true, "plan": true, "collection": true,
"hand-acts": true, "durations": true, "conditions": true, "doctor": true, "retire": true, "cleanup": true, "data": true, "hand-acts": true, "durations": true, "conditions": true, "doctor": true, "retire": true, "cleanup": true, "data": true,
// What the records say the registries may keep (novox/hq ADR 0251). // What the records say the registries may keep (novox/hq ADR 0251).
"artifacts": true, "collect": true, "images": true, "artifacts": true, "collect": true, "images": true, "mirrors": true,
// The delivery's owner's verbs answer JSON where they read (plan, order, check, walks) — novox/hq ADR 0239. // The delivery's owner's verbs answer JSON where they read (plan, order, check, walks) — novox/hq ADR 0239.
"delivery": true} "delivery": true}
@@ -835,6 +859,8 @@ func repairingCommand(argv []string) string {
return "cleanup delete" return "cleanup delete"
case argv[0] == "collect" && slices.Contains(argv, "--confirm"): case argv[0] == "collect" && slices.Contains(argv, "--confirm"):
return "collect" return "collect"
case argv[0] == "mirrors" && slices.Contains(argv, "--confirm"):
return "mirrors"
} }
return "" return ""
} }
@@ -284,6 +284,7 @@ var accountedFlags = map[string]map[string]string{
"artifacts": {"json": "set by the verb: the answer is data"}, "artifacts": {"json": "set by the verb: the answer is data"},
"collect": {"json": "set by the verb: the answer is data"}, "collect": {"json": "set by the verb: the answer is data"},
"images": {"json": "set by the verb: the answer is data"}, "images": {"json": "set by the verb: the answer is data"},
"mirrors": {"json": "set by the verb: the answer is data"},
"conditions history": {"json": "set by the verb: the answer is data"}, "conditions history": {"json": "set by the verb: the answer is data"},
"conditions show": {"json": "set by the verb: the answer is data"}, "conditions show": {"json": "set by the verb: the answer is data"},
"healers": {"json": "set by the verb: the answer is data"}, "healers": {"json": "set by the verb: the answer is data"},
+176
View File
@@ -0,0 +1,176 @@
package artifacts
import (
"context"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"slices"
"strings"
"sync"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// An index is let go of with the platform manifests it names, except those a kept index of the same
// repository names too (novox/hq ADR 0257). Each platform manifest is a manifest of the repository in
// its own right, and the store's collector keeps every manifest a repository holds — so an index let go
// of alone frees nothing of the images it names.
func digestN(n int) string { return fmt.Sprintf("sha256:%064x", n) }
// indexStore holds indexes and images in one repository, answers GETs with their documents, and records
// every delete.
type indexStore struct {
mu sync.Mutex
indexes map[string][]string // digest → the platform manifests it names
images map[string]bool
deleted []string
}
func (s *indexStore) serve(t *testing.T) Store {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
s.mu.Lock()
defer s.mu.Unlock()
digest := r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]
switch r.Method {
case http.MethodGet:
if children, ok := s.indexes[digest]; ok {
var named []string
for _, c := range children {
named = append(named, `{"mediaType":"application/vnd.oci.image.manifest.v1+json","digest":"`+c+`","size":1}`)
}
w.Header().Set("Content-Type", "application/vnd.oci.image.index.v1+json")
_, _ = w.Write([]byte(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.index.v1+json","manifests":[` +
strings.Join(named, ",") + `]}`))
return
}
if s.images[digest] {
w.Header().Set("Content-Type", "application/vnd.oci.image.manifest.v1+json")
_, _ = w.Write([]byte(`{"schemaVersion":2,"mediaType":"application/vnd.oci.image.manifest.v1+json","layers":[]}`))
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodHead:
if _, ok := s.indexes[digest]; ok || s.images[digest] {
w.WriteHeader(http.StatusOK)
return
}
w.WriteHeader(http.StatusNotFound)
case http.MethodDelete:
if _, ok := s.indexes[digest]; !ok && !s.images[digest] {
w.WriteHeader(http.StatusNotFound)
return
}
s.deleted = append(s.deleted, digest)
delete(s.indexes, digest)
delete(s.images, digest)
w.WriteHeader(http.StatusAccepted)
default:
w.WriteHeader(http.StatusBadRequest)
}
}))
t.Cleanup(server.Close)
return Store{Address: strings.TrimPrefix(server.URL, "http://")}
}
func TestAnIndexGoesWithItsPlatformsAndAKeptIndexKeepsItsOwn(t *testing.T) {
// Two indexes of one image's repository: the old one names platforms 11 and 12, the kept one 12 and 13.
s := &indexStore{
indexes: map[string][]string{digestN(1): {digestN(11), digestN(12)}, digestN(2): {digestN(12), digestN(13)}},
images: map[string]bool{digestN(11): true, digestN(12): true, digestN(13): true},
}
store := s.serve(t)
repository := "upstream/docker.io/library/golang"
old := catalogue.ArtifactStoreScheme + repository + "@" + digestN(1)
kept := catalogue.ArtifactStoreScheme + repository + "@" + digestN(2)
spare, err := store.SpareKeptIndexes(context.Background(), []string{kept}, []string{old})
if err != nil {
t.Fatal(err)
}
store.Spare = spare
if err := store.LetGo(context.Background(), old); err != nil {
t.Fatal(err)
}
// Its own platform first, the index last; the platform the kept index names is left in place.
if !slices.Equal(s.deleted, []string{digestN(11), digestN(1)}) {
t.Fatalf("deleted %v; want the old index's own platform, then the index", s.deleted)
}
if !s.images[digestN(12)] || !s.images[digestN(13)] {
t.Fatal("a platform a kept index names was let go of")
}
}
func TestWithoutSparingAnIndexGoesAloneAndAnImageNamesNoPlatforms(t *testing.T) {
s := &indexStore{indexes: map[string][]string{digestN(1): {digestN(11)}}, images: map[string]bool{digestN(11): true, digestN(5): true}}
store := s.serve(t)
// No Spare: as before ADR 0257, the index alone.
if err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/base@"+digestN(1)); err != nil {
t.Fatal(err)
}
if !slices.Equal(s.deleted, []string{digestN(1)}) {
t.Fatalf("without a Spare, deleted %v", s.deleted)
}
// An image is one delete, with a Spare or without.
s.deleted = nil
spare, err := store.SpareKeptIndexes(context.Background(), nil, []string{catalogue.ArtifactStoreScheme + "web/app@" + digestN(5)})
if err != nil {
t.Fatal(err)
}
store.Spare = spare
if err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/app@"+digestN(5)); err != nil {
t.Fatal(err)
}
if !slices.Equal(s.deleted, []string{digestN(5)}) {
t.Fatalf("an image was let go of as %v", s.deleted)
}
// An index the store no longer holds is Gone, and nothing is deleted on its account.
s.deleted = nil
err = store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/base@"+digestN(9))
if !errors.Is(err, Gone) || len(s.deleted) != 0 {
t.Fatalf("an index the store does not hold: %v, deleted %v", err, s.deleted)
}
}
func TestASpareThatCannotBeReadKeepsTheIndex(t *testing.T) {
s := &indexStore{indexes: map[string][]string{digestN(1): {digestN(11)}}, images: map[string]bool{digestN(11): true}}
store := s.serve(t)
store.Spare = func(context.Context, string) (map[string]bool, error) { return nil, fmt.Errorf("the store went away") }
err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/base@"+digestN(1))
if err == nil || !strings.Contains(err.Error(), "was kept") {
t.Fatalf("a spare that could not be read: %v", err)
}
if len(s.deleted) != 0 {
t.Fatalf("deleted %v without knowing what a kept index names", s.deleted)
}
}
// The spare list is read whole before the sweep: a kept index that cannot be read is an error before
// anything is deleted, and a repository not read before keeps its index.
func TestASpareListIsReadWholeBeforeAnythingGoes(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
}))
t.Cleanup(server.Close)
store := Store{Address: strings.TrimPrefix(server.URL, "http://")}
kept := catalogue.ArtifactStoreScheme + "web/base@" + digestN(2)
old := catalogue.ArtifactStoreScheme + "web/base@" + digestN(1)
if _, err := store.SpareKeptIndexes(context.Background(), []string{kept}, []string{old}); err == nil {
t.Fatal("a kept index the store would not answer for was taken as naming nothing")
}
s := &indexStore{indexes: map[string][]string{digestN(1): {digestN(11)}}, images: map[string]bool{digestN(11): true}}
good := s.serve(t)
spare, err := good.SpareKeptIndexes(context.Background(), nil, []string{catalogue.ArtifactStoreScheme + "other/repo@" + digestN(3)})
if err != nil {
t.Fatal(err)
}
good.Spare = spare
if err := good.LetGo(context.Background(), old); err == nil || len(s.deleted) != 0 {
t.Fatalf("an index in a repository not read before the sweep was let go of: %v, deleted %v", err, s.deleted)
}
}
+170
View File
@@ -9,8 +9,10 @@ package artifacts
import ( import (
"context" "context"
"encoding/json"
"errors" "errors"
"fmt" "fmt"
"io"
"net/http" "net/http"
"strings" "strings"
@@ -24,6 +26,10 @@ type Store struct {
Address string Address string
// HTTP is the client used; nil is a client with a modest timeout. // HTTP is the client used; nil is a client with a modest timeout.
HTTP *http.Client 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. // Gone is the answer when the store does not hold it: the outcome wanted, already true.
@@ -77,9 +83,142 @@ func (s Store) LetGo(ctx context.Context, reference string) error {
return err 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) 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. // 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 { func (s Store) remove(ctx context.Context, url, what string) error {
request, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil) request, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil)
@@ -121,3 +260,34 @@ func split(path string) (repository, kind, digest string, err error) {
} }
return "", "", "", fmt.Errorf("%w: %q names nothing the store holds by digest", ErrNotOurs, path) return "", "", "", fmt.Errorf("%w: %q names nothing the store holds by digest", ErrNotOurs, path)
} }
// IsIndex is whether a recorded reference is an index the store holds: one that names manifests. A
// blob, or a manifest that names none, is not. Gone when the store does not hold it.
func (s Store) IsIndex(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 || kind != "manifests" {
return false, err
}
children, err := s.Platforms(ctx, repository, digest)
return len(children) > 0, err
}
// 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...)
}
+42 -5
View File
@@ -77,6 +77,11 @@ type Result struct {
// the default, and that branch is its trunk. // the default, and that branch is its trunk.
Branches []string Branches []string
// Mirrored is every base this build copied into the artifact store, as pinned (novox/hq ADR 0257):
// said whether the build worked or not, because the copy is in the store either way, and the
// records are what decide whether it stays.
Mirrored []string
// Source is the build's source fingerprint (source.go): what it was made from — the module's tree, // Source is the build's source fingerprint (source.go): what it was made from — the module's tree,
// the contexts' trees, the bases and toolchains by digest — hashed. Empty where the source does not // the contexts' trees, the bases and toolchains by digest — hashed. Empty where the source does not
// pin the build. Two builds with one fingerprint are one build, whatever digests they made // pin the build. Two builds with one fingerprint are one build, whatever digests they made
@@ -106,6 +111,19 @@ type GitCredential struct {
func Build(ctx context.Context, run Runner, publish Publisher, func Build(ctx context.Context, run Runner, publish Publisher,
repository, path, ref, workspace string, held map[string]string, npmrc Npmrc, repository, path, ref, workspace string, held map[string]string, npmrc Npmrc,
forge GitCredential, log Log, seats ...map[string]string) (Result, error) { forge GitCredential, log Log, seats ...map[string]string) (Result, error) {
// **What was copied is said whether the build worked or not** (novox/hq ADR 0257): a base is in
// the store from the moment it is copied, and a build that failed after copying it is the
// one record that it is there.
var mirrored []string
result, err := build(ctx, run, publish, repository, path, ref, workspace, held, npmrc, forge, log,
&mirrored, seats...)
result.Mirrored = mirrored
return result, err
}
func build(ctx context.Context, run Runner, publish Publisher,
repository, path, ref, workspace string, held map[string]string, npmrc Npmrc,
forge GitCredential, log Log, mirrored *[]string, seats ...map[string]string) (Result, error) {
// The clone base of each seat a context may name (novox/hq ADR 0155); variadic so the callers // The clone base of each seat a context may name (novox/hq ADR 0155); variadic so the callers
// that hand none — tests of everything but contexts — read as they did. // that hand none — tests of everything but contexts — read as they did.
var seatBases map[string]string var seatBases map[string]string
@@ -222,10 +240,22 @@ func Build(ctx context.Context, run Runner, publish Publisher,
// An image published elsewhere that the build stands on is copied into the mesh's own // An image published elsewhere that the build stands on is copied into the mesh's own
// registry first, like an upstream artifact (ADR 0096), and the recipe is handed the copy. // registry first, like an upstream artifact (ADR 0096), and the recipe is handed the copy.
// Genesis has nowhere to copy to and pulls it into this machine's store instead. // Genesis has nowhere to copy to and pulls it into this machine's store instead.
mirror := func(ctx context.Context, from, repository string) (string, error) { mirror := func(ctx context.Context, from, repository, former string) (string, error) {
if m, can := publish.(BaseMirrorer); can {
say("bases", "copying %s into the mesh's registry as %s", from, repository)
reference, err := m.MirrorBase(ctx, from, repository, former)
if err == nil {
*mirrored = append(*mirrored, reference)
}
return reference, err
}
if m, can := publish.(Mirrorer); can { if m, can := publish.(Mirrorer); can {
say("bases", "copying %s into the mesh's registry", from) say("bases", "copying %s into the mesh's registry as %s", from, repository)
return m.MirrorImage(ctx, from, repository) reference, err := m.MirrorImage(ctx, from, repository)
if err == nil {
*mirrored = append(*mirrored, reference)
}
return reference, err
} }
if _, err := run(ctx, tree, "docker", "pull", from); err != nil { if _, err := run(ctx, tree, "docker", "pull", from); err != nil {
return "", fmt.Errorf("cannot fetch %s: %w", from, err) return "", fmt.Errorf("cannot fetch %s: %w", from, err)
@@ -901,7 +931,7 @@ var _ io.Writer = (*stringWriter)(nil)
// The order is fixed so two builds of one commit invoke the same command. Returned alongside the // The order is fixed so two builds of one commit invoke the same command. Returned alongside the
// arguments is every reference they resolved to, which is what the build stood on. // arguments is every reference they resolved to, which is what the build stood on.
func standingOn(ctx context.Context, manifest catalogue.Manifest, held map[string]string, func standingOn(ctx context.Context, manifest catalogue.Manifest, held map[string]string,
mirror func(ctx context.Context, from, repository string) (string, error)) ([]string, []string, error) { mirror func(ctx context.Context, from, repository, former string) (string, error)) ([]string, []string, error) {
if manifest.Build == nil || len(manifest.Build.On) == 0 { if manifest.Build == nil || len(manifest.Build.On) == 0 {
return nil, nil, nil return nil, nil, nil
} }
@@ -924,7 +954,14 @@ func standingOn(ctx context.Context, manifest catalogue.Manifest, held map[strin
"%s stands on the image %q, which is not pinned by digest. A tag is what "+ "%s stands on the image %q, which is not pinned by digest. A tag is what "+
"somebody else can move; name it as <image>@sha256:…", manifest.Module, base.Image) "somebody else can move; name it as <image>@sha256:…", manifest.Module, base.Image)
} }
reference, err := mirror(ctx, base.Image, manifest.Module+"/on-"+strings.ToLower(base.Arg)) // **One copy per upstream image, whichever modules stand on it** (novox/hq ADR 0257). Its
// repository is named for the image, not the module; the module's own repository of
// before is offered as a source, so the move to one copy asks upstream for nothing.
repository, err := MirrorRepository(base.Image)
if err != nil {
return nil, nil, fmt.Errorf("%s stands on %s: %w", manifest.Module, base.Image, err)
}
reference, err := mirror(ctx, base.Image, repository, FormerMirrorRepository(manifest.Module, base.Arg))
if err != nil { if err != nil {
return nil, nil, fmt.Errorf("%s stands on %s: %w", manifest.Module, base.Image, err) return nil, nil, fmt.Errorf("%s stands on %s: %w", manifest.Module, base.Image, err)
} }
+117 -10
View File
@@ -27,6 +27,40 @@ type Mirrorer interface {
MirrorImage(ctx context.Context, from, repository string) (string, error) MirrorImage(ctx context.Context, from, repository string) (string, error)
} }
// BaseMirrorer copies a base a build stands on into the one repository the mesh keeps for that
// upstream image, taking what an earlier copy under `former` already holds rather than asking
// upstream for it again (novox/hq ADR 0257).
type BaseMirrorer interface {
MirrorBase(ctx context.Context, from, repository, former string) (string, error)
}
// MirrorPrefix is the namespace of the repositories a base is mirrored into: `upstream/<host>/<path>`,
// one repository per upstream image, whichever modules stand on it (novox/hq ADR 0257).
const MirrorPrefix = "upstream/"
// MirrorRepository is the repository the mesh keeps an upstream image's copy in: the image's own
// host and path under MirrorPrefix, so two modules standing on one image stand on one copy, and two
// images that share a path on different hosts are never confused. Lower case and without a port's
// colon, because a repository name allows neither.
func MirrorRepository(from string) (string, error) {
where, err := parseReference(from)
if err != nil {
return "", err
}
host := strings.TrimPrefix(strings.TrimPrefix(where.base, "https://"), "http://")
if host == "registry-1.docker.io" {
host = "docker.io"
}
host = strings.ReplaceAll(host, ":", "-")
return strings.ToLower(MirrorPrefix + host + "/" + where.repository), nil
}
// FormerMirrorRepository is where a module's base was copied before ADR 0257: under the module's own
// repository, one copy per module (ADR 0097). Read only, as a source of what is already held.
func FormerMirrorRepository(module, arg string) string {
return module + "/on-" + strings.ToLower(arg)
}
const ( const (
mediaIndexOCI = "application/vnd.oci.image.index.v1+json" mediaIndexOCI = "application/vnd.oci.image.index.v1+json"
mediaIndexDocker = "application/vnd.docker.distribution.manifest.list.v2+json" mediaIndexDocker = "application/vnd.docker.distribution.manifest.list.v2+json"
@@ -86,6 +120,9 @@ func parseReference(ref string) (upstream, error) {
type source struct { type source struct {
client *http.Client client *http.Client
token string token string
// mountFrom is a repository of the mesh's own registry the source is, so a blob is mounted from
// it rather than read and written again.
mountFrom string
} }
// get fetches a registry URL, answering a bearer challenge once with an anonymous token — which is // get fetches a registry URL, answering a bearer challenge once with an anonymous token — which is
@@ -176,6 +213,15 @@ type descriptor struct {
// under `repository`, and returns the reference the mesh will pin: this registry, the repository, // under `repository`, and returns the reference the mesh will pin: this registry, the repository,
// and the digest of the document that was put last, which is the index where there is one. // and the digest of the document that was put last, which is the index where there is one.
func (r Registry) MirrorImage(ctx context.Context, from, repository string) (string, error) { func (r Registry) MirrorImage(ctx context.Context, from, repository string) (string, error) {
return r.MirrorBase(ctx, from, repository, "")
}
// MirrorBase is MirrorImage, and where `former` — a repository of this registry — already holds the
// image by its digest, the copy is made from there: manifests read from this registry, blobs mounted
// across rather than moved (novox/hq ADR 0257). Upstream is asked only for what no repository here
// holds, which is what kept the public hub's anonymous pull limit out of reach when every module's
// base moved into one shared repository.
func (r Registry) MirrorBase(ctx context.Context, from, repository, former string) (string, error) {
where, err := parseReference(from) where, err := parseReference(from)
if err != nil { if err != nil {
return "", err return "", err
@@ -186,7 +232,11 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
// on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base // on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base
// lives there failed on a copy it did not need. // lives there failed on a copy it did not need.
if strings.HasPrefix(where.reference, "sha256:") { if strings.HasPrefix(where.reference, "sha256:") {
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference, manifestAccept) // **Held whole, not only its index** (novox/hq ADR 0257): every platform manifest the index
// names is asked about too. A sweep that stopped half-way, or that ran while this build was
// copying, may have left an index naming a platform the store no longer holds; such a copy
// is made again, and the copy puts back what is missing.
held, err := r.holdsWhole(ctx, repository, where.reference)
if err != nil { if err != nil {
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err) return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
} }
@@ -195,6 +245,17 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
} }
} }
src := &source{client: r.client()} src := &source{client: r.client()}
if former != "" && former != repository && strings.HasPrefix(where.reference, "sha256:") {
held, err := r.holdsWhole(ctx, former, where.reference)
if err != nil {
return "", fmt.Errorf("asking %s whether %s holds %s: %w", r.Address, former, from, err)
}
if held {
// The same bytes, already here: copied from this registry, each blob mounted.
where = upstream{base: "http://" + r.Address, repository: former, reference: where.reference}
src.mountFrom = former
}
}
digest, err := r.copyManifest(ctx, src, where, where.reference, repository) digest, err := r.copyManifest(ctx, src, where, where.reference, repository)
if err != nil { if err != nil {
return "", fmt.Errorf("copying %s into %s/%s: %w", from, r.Address, repository, err) return "", fmt.Errorf("copying %s into %s/%s: %w", from, r.Address, repository, err)
@@ -202,6 +263,42 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
return r.Address + "/" + repository + "@" + digest, nil return r.Address + "/" + repository + "@" + digest, nil
} }
// holdsWhole is whether this registry holds the manifest under the repository and, for an index, every
// manifest it names.
func (r Registry) holdsWhole(ctx context.Context, repository, digest string) (bool, error) {
url := "http://" + r.Address + "/v2/" + repository + "/manifests/"
held, err := r.has(ctx, url+digest, manifestAccept)
if err != nil || !held {
return false, err
}
request, err := http.NewRequestWithContext(ctx, http.MethodGet, url+digest, nil)
if err != nil {
return false, err
}
request.Header.Set("Accept", manifestAccept)
response, err := r.client().Do(request)
if err != nil {
return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err)
}
defer response.Body.Close()
if response.StatusCode != http.StatusOK {
return false, nil
}
var document struct {
Manifests []descriptor `json:"manifests"`
}
if err := json.NewDecoder(io.LimitReader(response.Body, 4<<20)).Decode(&document); err != nil {
return false, fmt.Errorf("%s@%s is not a manifest: %w", repository, digest, err)
}
for _, m := range document.Manifests {
held, err := r.has(ctx, url+m.Digest, manifestAccept)
if err != nil || !held {
return false, err
}
}
return true, nil
}
// copyManifest copies one manifest document and everything it names, and returns its digest. An // copyManifest copies one manifest document and everything it names, and returns its digest. An
// index is copied by copying each manifest it names first, so the index never points at something // index is copied by copying each manifest it names first, so the index never points at something
// the registry does not hold yet. // the registry does not hold yet.
@@ -286,16 +383,14 @@ func (r Registry) copyBlob(ctx context.Context, src *source, where upstream, dig
} else if there { } else if there {
return nil return nil
} }
response, err := src.get(ctx, where.base+"/v2/"+where.repository+"/blobs/"+digest, "") // **Mounted where this registry already holds it** (novox/hq ADR 0257): a blob is stored once
if err != nil { // whichever repositories link it, so a mount moves no bytes. A registry that cannot mount answers
return err // with an ordinary upload's location, and the blob is moved as before.
uploads := base + "/blobs/uploads/"
if src.mountFrom != "" {
uploads += "?mount=" + digest + "&from=" + src.mountFrom
} }
defer response.Body.Close() start, err := http.NewRequestWithContext(ctx, http.MethodPost, uploads, nil)
if response.StatusCode != http.StatusOK {
return fmt.Errorf("%s/%s: blob %s: %s", where.base, where.repository, digest, response.Status)
}
start, err := http.NewRequestWithContext(ctx, http.MethodPost, base+"/blobs/uploads/", nil)
if err != nil { if err != nil {
return err return err
} }
@@ -304,9 +399,21 @@ func (r Registry) copyBlob(ctx context.Context, src *source, where upstream, dig
return fmt.Errorf("cannot start an upload to %s: %w", base, err) return fmt.Errorf("cannot start an upload to %s: %w", base, err)
} }
begun.Body.Close() begun.Body.Close()
if begun.StatusCode == http.StatusCreated && src.mountFrom != "" {
return nil
}
if begun.StatusCode != http.StatusAccepted { if begun.StatusCode != http.StatusAccepted {
return fmt.Errorf("%s answered %s when asked where to put a blob", base, begun.Status) return fmt.Errorf("%s answered %s when asked where to put a blob", base, begun.Status)
} }
response, err := src.get(ctx, where.base+"/v2/"+where.repository+"/blobs/"+digest, "")
if err != nil {
return err
}
defer response.Body.Close()
if response.StatusCode != http.StatusOK {
return fmt.Errorf("%s/%s: blob %s: %s", where.base, where.repository, digest, response.Status)
}
location := begun.Header.Get("Location") location := begun.Header.Get("Location")
if location == "" { if location == "" {
return fmt.Errorf("%s accepted an upload and said nowhere to put it", base) return fmt.Errorf("%s accepted an upload and said nowhere to put it", base)
+375
View File
@@ -0,0 +1,375 @@
package builder
import (
"context"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
)
// One copy per upstream image, whichever modules stand on it (novox/hq ADR 0257).
func TestAnUpstreamImageHasOneRepositoryNamedForIt(t *testing.T) {
for from, want := range map[string]string{
"golang@sha256:abc": "upstream/docker.io/library/golang",
"n8nio/n8n@sha256:abc": "upstream/docker.io/n8nio/n8n",
"quay.io/minio/mc@sha256:abc": "upstream/quay.io/minio/mc",
"ghcr.io/Mailu/Admin@sha256:abc": "upstream/ghcr.io/mailu/admin",
"localhost:5000/x/y@sha256:abc": "upstream/localhost-5000/x/y",
"docker.io/library/alpine:3.20@sha256:ab": "upstream/docker.io/library/alpine",
} {
got, err := MirrorRepository(from)
if err != nil {
t.Fatalf("%s: %v", from, err)
}
if got != want {
t.Errorf("%s: got %q want %q", from, got, want)
}
}
if got := FormerMirrorRepository("route-proxy", "GO_BASE"); got != "route-proxy/on-go_base" {
t.Fatalf("the former repository is %q", got)
}
}
// aRegistryOfRepositories is a registry the way the real one is: blobs stored once, and each
// repository linking the blobs and manifests it holds. It mounts a blob across repositories.
type aRegistryOfRepositories struct {
mu sync.Mutex
blobs map[string][]byte // digest → bytes, stored once
links map[string]map[string]bool // repository → blob digests it links
manifests map[string][]byte // repository@digest → document
uploads int
mounts int
gets int
// noMount answers every mount with an ordinary upload's location, as a registry that cannot
// mount across repositories does.
noMount bool
}
func newRegistryOfRepositories() *aRegistryOfRepositories {
return &aRegistryOfRepositories{blobs: map[string][]byte{}, links: map[string]map[string]bool{},
manifests: map[string][]byte{}}
}
func (m *aRegistryOfRepositories) link(repository, digest string) {
if m.links[repository] == nil {
m.links[repository] = map[string]bool{}
}
m.links[repository][digest] = true
}
// split reads /v2/<repository>/<kind>/<rest> with a repository of any depth.
func splitPath(path string) (repository, kind, rest string) {
path = strings.TrimPrefix(path, "/v2/")
for _, k := range []string{"/manifests/", "/blobs/uploads/", "/blobs/"} {
if i := strings.Index(path, k); i >= 0 {
return path[:i], strings.Trim(k, "/"), path[i+len(k):]
}
}
return "", "", ""
}
func (m *aRegistryOfRepositories) handler() http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
m.mu.Lock()
defer m.mu.Unlock()
repository, kind, rest := splitPath(r.URL.Path)
switch {
case kind == "manifests" && (r.Method == http.MethodHead || r.Method == http.MethodGet):
body, ok := m.manifests[repository+"@"+rest]
if !ok {
w.WriteHeader(http.StatusNotFound)
return
}
if r.Method == http.MethodGet {
m.gets++
w.Header().Set("Content-Type", mediaTypeOf(body))
_, _ = w.Write(body)
return
}
w.WriteHeader(http.StatusOK)
case kind == "manifests" && r.Method == http.MethodPut:
body, _ := readAll(r)
m.manifests[repository+"@"+digestOf(body)] = body
w.WriteHeader(http.StatusCreated)
case kind == "blobs" && (r.Method == http.MethodHead || r.Method == http.MethodGet):
if !m.links[repository][rest] {
w.WriteHeader(http.StatusNotFound)
return
}
if r.Method == http.MethodGet {
m.gets++
_, _ = w.Write(m.blobs[rest])
return
}
w.WriteHeader(http.StatusOK)
case kind == "manifests" && r.Method == http.MethodDelete:
delete(m.manifests, repository+"@"+rest)
w.WriteHeader(http.StatusAccepted)
case kind == "blobs/uploads" && r.Method == http.MethodPost:
if digest, from := r.URL.Query().Get("mount"), r.URL.Query().Get("from"); digest != "" && m.links[from][digest] && !m.noMount {
m.link(repository, digest)
m.mounts++
w.WriteHeader(http.StatusCreated)
return
}
w.Header().Set("Location", "/v2/"+repository+"/blobs/uploads/one")
w.WriteHeader(http.StatusAccepted)
case kind == "blobs/uploads" && r.Method == http.MethodPut:
body, _ := readAll(r)
digest := r.URL.Query().Get("digest")
if digestOf(body) != digest {
w.WriteHeader(http.StatusBadRequest)
return
}
m.blobs[digest] = body
m.link(repository, digest)
m.uploads++
w.WriteHeader(http.StatusCreated)
default:
w.WriteHeader(http.StatusNotFound)
}
})
}
func mediaTypeOf(body []byte) string {
switch {
case strings.Contains(string(body), mediaIndexOCI):
return mediaIndexOCI
default:
return mediaManifestOCI
}
}
// A base a module's own repository already holds is copied into the image's repository from there:
// every blob mounted, nothing asked of upstream — which may be gone, or rationing anonymous pulls.
func TestABaseHeldUnderTheModulesFormerRepositoryIsMountedNotFetchedAgain(t *testing.T) {
src, indexDigest, srcBlobs := anUpstreamRegistry(t)
dst := newRegistryOfRepositories()
dstServer := httptest.NewServer(dst.handler())
defer dstServer.Close()
address := strings.TrimPrefix(dstServer.URL, "http://")
r := Registry{Address: address, HTTP: src.Client()}
host := strings.TrimPrefix(src.URL, "http://")
// Before: copied the old way, under the module's repository.
if _, err := r.MirrorImage(context.Background(), host+"/library/thing:latest", "hello-web/on-thing_base"); err != nil {
t.Fatal(err)
}
uploaded := dst.uploads
if uploaded != len(srcBlobs) {
t.Fatalf("the first copy uploaded %d of %d blobs", uploaded, len(srcBlobs))
}
src.Close()
from := host + "/library/thing@" + indexDigest
repository, err := MirrorRepository(from)
if err != nil {
t.Fatal(err)
}
reference, err := r.MirrorBase(context.Background(), from, repository, "hello-web/on-thing_base")
if err != nil {
t.Fatalf("a base the registry holds was asked of an upstream that is gone: %v", err)
}
if reference != address+"/"+repository+"@"+indexDigest {
t.Fatalf("pinned as %q", reference)
}
if dst.uploads != uploaded {
t.Fatalf("the move to one repository uploaded %d blobs again", dst.uploads-uploaded)
}
if dst.mounts != len(srcBlobs) {
t.Fatalf("%d of %d blobs mounted", dst.mounts, len(srcBlobs))
}
// The index and both platform manifests are in the image's repository, and every blob is linked.
for key := range dst.manifests {
if strings.HasPrefix(key, "hello-web/") {
continue
}
if !strings.HasPrefix(key, repository+"@") {
t.Fatalf("a manifest went to %s", key)
}
}
if n := countPrefix(dst.manifests, repository+"@"); n != 3 {
t.Fatalf("%d manifests in %s, want the index and two platforms", n, repository)
}
for digest := range srcBlobs {
if !dst.links[repository][digest] {
t.Fatalf("blob %s is not linked into %s", digest, repository)
}
}
// A second module standing on the same image: already held, nothing copied, the same reference.
again, err := r.MirrorBase(context.Background(), from, repository, "other-module/on-thing_base")
if err != nil || again != reference {
t.Fatalf("a second module's copy: %q %v", again, err)
}
}
func countPrefix(m map[string][]byte, prefix string) int {
n := 0
for k := range m {
if strings.HasPrefix(k, prefix) {
n++
}
}
return n
}
// What a build copied is said whether it worked or not: the copy is in the store either way, and a
// build that failed after copying is the one record that it is there.
func TestABuildSaysWhatItMirroredEvenWhenItFails(t *testing.T) {
// A recipe that reaches for an image nobody declared is refused — after the base was copied.
r, workspace := aRepository(t, anImage, map[string]string{
"modules/bus/Dockerfile": "ARG BASE\nFROM ${BASE}\nCOPY --from=vendor/other:1 /x /x"})
r.contents["modules/bus/module.json"] = anImage
r.tree = "aaaa"
m := &mirroring{recorded: r, at: "registry-a:5000"}
got, err := Build(context.Background(), r.run, m, "https://forge.invalid/catalogue.git", "modules/bus", "", workspace,
nil, Npmrc{}, GitCredential{}, nil)
if err == nil {
t.Fatal("a recipe fetching an undeclared image was built")
}
want := "registry-a:5000/upstream/docker.io/vendor/server@sha256:" + strings.Repeat("1", 64)
if len(got.Mirrored) != 1 || got.Mirrored[0] != want {
t.Fatalf("a failed build said it mirrored %v, want [%s]", got.Mirrored, want)
}
// And a build that works says it too.
r2, workspace2 := aRepository(t, anImage, map[string]string{"modules/bus/Dockerfile": "ARG BASE\nFROM ${BASE}"})
r2.contents["modules/bus/module.json"] = anImage
r2.tree = "aaaa"
ok, err := Build(context.Background(), r2.run, &mirroring{recorded: r2, at: "registry-a:5000"},
"https://forge.invalid/catalogue.git", "modules/bus", "", workspace2, nil, Npmrc{}, GitCredential{}, nil)
if err != nil {
t.Fatal(err)
}
if len(ok.Mirrored) != 1 || ok.Mirrored[0] != want {
t.Fatalf("a build said it mirrored %v, want [%s]", ok.Mirrored, want)
}
}
// held puts one image in a repository the old way and answers its index digest and the registry.
func heldUnderAModule(t *testing.T) (*aRegistryOfRepositories, Registry, string, string, map[string][]byte, func()) {
t.Helper()
src, indexDigest, srcBlobs := anUpstreamRegistry(t)
t.Cleanup(src.Close)
dst := newRegistryOfRepositories()
dstServer := httptest.NewServer(dst.handler())
t.Cleanup(dstServer.Close)
r := Registry{Address: strings.TrimPrefix(dstServer.URL, "http://"), HTTP: src.Client()}
host := strings.TrimPrefix(src.URL, "http://")
if _, err := r.MirrorImage(context.Background(), host+"/library/thing:latest", "hello-web/on-thing_base"); err != nil {
t.Fatal(err)
}
return dst, r, host + "/library/thing@" + indexDigest, indexDigest, srcBlobs, src.Close
}
// An index whose platform manifests are not all held is not "already held": it is copied again, and the
// copy puts back what is missing — what a sweep that stopped half-way, or ran while a build was copying,
// leaves behind.
func TestAnIndexMissingAPlatformIsCopiedAgain(t *testing.T) {
dst, r, from, indexDigest, _, _ := heldUnderAModule(t)
repository, _ := MirrorRepository(from)
if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil {
t.Fatal(err)
}
// A sweep lets one platform go, between this build's copy and the next.
var platform string
for key := range dst.manifests {
if strings.HasPrefix(key, repository+"@") && !strings.HasSuffix(key, indexDigest) {
platform = key
break
}
}
delete(dst.manifests, platform)
if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil {
t.Fatal(err)
}
if _, back := dst.manifests[platform]; !back {
t.Fatalf("an index missing %s was taken as held", platform)
}
}
// The same, for the module's former copy: a former copy missing a platform is not a source; upstream
// is asked instead.
func TestAFormerCopyMissingAPlatformIsNotTheSource(t *testing.T) {
dst, r, from, indexDigest, _, _ := heldUnderAModule(t)
for key := range dst.manifests {
if strings.HasPrefix(key, "hello-web/on-thing_base@") && !strings.HasSuffix(key, indexDigest) {
delete(dst.manifests, key)
break
}
}
repository, _ := MirrorRepository(from)
if _, err := r.MirrorBase(context.Background(), from, repository, "hello-web/on-thing_base"); err != nil {
t.Fatal(err)
}
if dst.mounts != 0 {
t.Fatalf("a former copy missing a platform was mounted from (%d mounts)", dst.mounts)
}
if n := countPrefix(dst.manifests, repository+"@"); n != 3 {
t.Fatalf("%d manifests copied from upstream, want 3", n)
}
}
// A registry that answers a mount with an upload's location (202) gets the blob moved as before.
func TestAMountRefusedFallsBackToAnUpload(t *testing.T) {
dst, r, from, _, srcBlobs, _ := heldUnderAModule(t)
dst.noMount = true
uploaded := dst.uploads
repository, _ := MirrorRepository(from)
reference, err := r.MirrorBase(context.Background(), from, repository, "hello-web/on-thing_base")
if err != nil {
t.Fatal(err)
}
if !strings.HasSuffix(reference, repository+"@"+strings.SplitN(from, "@", 2)[1]) {
t.Fatalf("pinned as %s", reference)
}
if dst.mounts != 0 || dst.uploads-uploaded != len(srcBlobs) {
t.Fatalf("%d mounts, %d uploads; want every blob uploaded from the former copy", dst.mounts, dst.uploads-uploaded)
}
for digest := range srcBlobs {
if !dst.links[repository][digest] {
t.Fatalf("blob %s is not linked into %s", digest, repository)
}
}
}
// A sweep and a build at once: while builds copy the base, platforms are let go of under them. Each
// build that returns has the copy whole.
func TestABuildCopyingWhileASweepLetsGoEndsWhole(t *testing.T) {
dst, r, from, indexDigest, _, _ := heldUnderAModule(t)
repository, _ := MirrorRepository(from)
if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil {
t.Fatal(err)
}
done := make(chan struct{})
go func() {
defer close(done)
for i := 0; i < 50; i++ {
dst.mu.Lock()
for key := range dst.manifests {
if strings.HasPrefix(key, repository+"@") && !strings.HasSuffix(key, indexDigest) {
delete(dst.manifests, key)
break
}
}
dst.mu.Unlock()
}
}()
for i := 0; i < 20; i++ {
if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil {
t.Fatal(err)
}
}
<-done
// The sweep is over; the next build finds what it left and makes the copy whole.
if _, err := r.MirrorBase(context.Background(), from, repository, ""); err != nil {
t.Fatal(err)
}
if n := countPrefix(dst.manifests, repository+"@"); n != 3 {
t.Fatalf("after the sweep and the builds, %d manifests in %s, want 3", n, repository)
}
}
+7
View File
@@ -117,6 +117,13 @@ func (m *theMeshsRegistry) handler() http.Handler {
} else { } else {
w.WriteHeader(http.StatusNotFound) w.WriteHeader(http.StatusNotFound)
} }
case r.Method == http.MethodGet && strings.Contains(r.URL.Path, "/manifests/"):
body, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]
if !ok {
w.WriteHeader(http.StatusNotFound)
return
}
_, _ = w.Write(body)
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"): case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"):
if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok { if _, ok := m.blobs[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
+9 -6
View File
@@ -72,7 +72,7 @@ func TestAnIncompleteBaseIsRefused(t *testing.T) {
} }
// noMirror is a mirror for tests whose bases are all the mesh's own. // noMirror is a mirror for tests whose bases are all the mesh's own.
func noMirror(context.Context, string, string) (string, error) { func noMirror(context.Context, string, string, string) (string, error) {
return "", fmt.Errorf("nothing to copy in this test") return "", fmt.Errorf("nothing to copy in this test")
} }
@@ -86,17 +86,20 @@ func TestADeclaredVendorImageIsCopiedInAndHandedToTheRecipe(t *testing.T) {
}, },
} }
var asked []string var asked []string
args, _, err := standingOn(context.Background(), manifest, nil, func(_ context.Context, from, repository string) (string, error) { args, _, err := standingOn(context.Background(), manifest, nil, func(_ context.Context, from, repository, former string) (string, error) {
asked = append(asked, from+" -> "+repository) asked = append(asked, from+" -> "+repository+" (from "+former+")")
return "127.0.0.1:5000/" + repository + "@sha256:" + strings.Repeat("d", 64), nil return "127.0.0.1:5000/" + repository + "@sha256:" + strings.Repeat("d", 64), nil
}) })
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
} }
if len(asked) != 1 || asked[0] != "quay.io/minio/mc@sha256:"+strings.Repeat("c", 64)+" -> minio/on-mc_base" { // One repository per upstream image, named for the image (novox/hq ADR 0257); the module's
t.Fatalf("the image was not copied under the module's repository: %v", asked) // repository of before offered as where it may already be held.
if len(asked) != 1 || asked[0] != "quay.io/minio/mc@sha256:"+strings.Repeat("c", 64)+
" -> upstream/quay.io/minio/mc (from minio/on-mc_base)" {
t.Fatalf("the image was not copied into the image's own repository: %v", asked)
} }
if strings.Join(args, " ") != "--build-arg MC_BASE=127.0.0.1:5000/minio/on-mc_base@sha256:"+strings.Repeat("d", 64) { if strings.Join(args, " ") != "--build-arg MC_BASE=127.0.0.1:5000/upstream/quay.io/minio/mc@sha256:"+strings.Repeat("d", 64) {
t.Fatalf("the recipe was not handed the copy: %v", args) t.Fatalf("the recipe was not handed the copy: %v", args)
} }
// Unpinned, it is refused: a tag is what somebody else can move. // Unpinned, it is refused: a tag is what somebody else can move.
+18 -1
View File
@@ -449,7 +449,10 @@ var ControllerVerbs = []Verb{
{Name: "collect", Description: "The artifact store's sweep, now: what the mesh made and keeps for no reason " + {Name: "collect", Description: "The artifact store's sweep, now: what the mesh made and keeps for no reason " +
"is let go of. Without confirm a dry run — every kept archive asked about with HEADs, what would be let go " + "is let go of. Without confirm a dry run — every kept archive asked about with HEADs, what would be let go " +
"listed, nothing held or deleted. With confirm and why: every kept archive held first, then each eligible " + "listed, nothing held or deleted. With confirm and why: every kept archive held first, then each eligible " +
"artifact let go of, oldest first, and recorded collected — a hand act, recorded with its why. Bounded by " + "artifact let go of, oldest first, and recorded collected — a hand act, recorded with its why. A confirmed " +
"collect lets an index go with the platform manifests no kept index names (read first; any read that " +
"fails lets nothing go); the sweep after a build leaves an eligible index for such a collect (novox/hq ADR " +
"0257). Bounded by " +
"most (500 by default, at most 5000) and 45 seconds; what is left is said. Bytes are reclaimed by the " + "most (500 by default, at most 5000) and 45 seconds; what is left is said. Bytes are reclaimed by the " +
"store's nightly collector. Never touches a digest the mesh did not record (novox/hq ADR 0189, ADR 0251).", "store's nightly collector. Never touches a digest the mesh did not record (novox/hq ADR 0189, ADR 0251).",
Input: schema(map[string]string{ Input: schema(map[string]string{
@@ -457,6 +460,20 @@ var ControllerVerbs = []Verb{
"why": "with confirm: why — required, and recorded in the hand-act log", "why": "with confirm: why — required, and recorded in the hand-act log",
"most": "let go of at most this many (500 by default, at most 5000)", "most": "let go of at most this many (500 by default, at most 5000)",
}, nil, "confirm")}, }, nil, "confirm")},
{Name: "mirrors", Description: "The bases the mesh copied into its artifact store for builds to stand on, as " +
"the records hold them: each with its state — kept because a kept build stood on it, or eligible — and " +
"which build or person recorded it. With record: references of copies no record names (a copy made " +
"before copies were recorded, or by a build that failed before then) to record as the mesh's, so the " +
"sweep keeps or lets go of them like any other; only a repository the mirror writes is taken " +
"(upstream/<host>/<path>, or a module's <module>/on-<argument>), and each must be in the store. Without " +
"confirm a dry run; with confirm and why, recorded — a hand act. A recorded copy is eligible: the next " +
"sweep lets it go unless a kept build stood on it. Records only: this verb deletes nothing " +
"(novox/hq ADR 0257).",
Input: schema(map[string]string{
"record": "references to record, artifact-store://<repository>@sha256:<hex>, separated by spaces or commas",
"confirm": "\"true\": record them (needs why); without it nothing is recorded",
"why": "with confirm: why — required, and recorded with each copy and in the hand-act log",
}, nil, "confirm")},
{Name: "images", Description: "Every container image one machine's declaration names — the declaration the " + {Name: "images", Description: "Every container image one machine's declaration names — the declaration the " +
"mesh would send it now and the one it was last sent — each with the resources naming it: what the " + "mesh would send it now and the one it was last sent — each with the resources naming it: what the " +
"machine keeps when it prunes its images. sent_known is false when the last one sent is not kept with " + "machine keeps when it prunes its images. sent_known is false when the last one sent is not kept with " +
+38 -1
View File
@@ -39,6 +39,9 @@ type Build struct {
// Against is every artifact this build stood on, as references rather than module names — // Against is every artifact this build stood on, as references rather than module names —
// what makes a build edge derived rather than declared (ADR 0009). // what makes a build edge derived rather than declared (ADR 0009).
Against []string Against []string
// Mirrored is every base this build copied into the artifact store (novox/hq ADR 0257), said
// whether the build worked or not: the copy is the mesh's, recorded in artifact_mirrored.
Mirrored []string
// Read is every repository this build read source from besides the module's own (novox/hq // Read is every repository this build read source from besides the module's own (novox/hq
// 04-ISSUES/131), at the ref it read. // 04-ISSUES/131), at the ref it read.
Read []ReadRepository Read []ReadRepository
@@ -120,7 +123,41 @@ func (i *Inventory) RecordBuild(ctx context.Context, b Build) error {
on conflict (id) do nothing`, on conflict (id) do nothing`,
b.ID, b.Repository, b.Ref, module, b.Commit, b.On, b.Failed, made, b.ID, b.Repository, b.Ref, module, b.Commit, b.On, b.Failed, made,
b.Path, manifestOrNil(b.Manifest), against, read, asked, b.SourceFingerprint) b.Path, manifestOrNil(b.Manifest), against, read, asked, b.SourceFingerprint)
return err if err != nil {
return err
}
return i.RecordMirrored(ctx, b.ID, "", b.Mirrored)
}
// RecordMirrored keeps that the mesh copied these bases into the artifact store (novox/hq ADR 0257):
// by the build that said so, or by a person who says why. The first record of a copy stands.
//
// **A copy a build says it holds is held again.** The builder answers a copy only once the store holds
// it whole, copying again what was let go of; so a copy the records say was collected, and a build now
// says it copied, is no longer collected. Left collected, it would stay in the store and out of every
// sweep for ever.
func (i *Inventory) RecordMirrored(ctx context.Context, build, why string, references []string) error {
var by *string
if build != "" {
by = &build
}
for _, reference := range references {
if reference == "" {
continue
}
if _, err := i.store.Pool().Exec(ctx,
`insert into artifact_mirrored (reference, build, why) values ($1, $2, $3)
on conflict (reference) do nothing`, reference, by, why); err != nil {
return err
}
if build != "" {
if _, err := i.store.Pool().Exec(ctx,
`delete from artifact_collected where reference = $1`, reference); err != nil {
return err
}
}
}
return nil
} }
// Builds is what has happened lately, newest first. // Builds is what has happened lately, newest first.
+118 -10
View File
@@ -37,6 +37,10 @@ const (
// KeptUnaddressable: it names no digest, so nothing here can speak for it, and it is kept rather // KeptUnaddressable: it names no digest, so nothing here can speak for it, and it is kept rather
// than guessed about. // than guessed about.
KeptUnaddressable = "unaddressable" KeptUnaddressable = "unaddressable"
// KeptStoodOn: it is a base the mesh copied in, and one of the KeptBuilds most recent successful
// builds of a module the mesh still holds stood on it (novox/hq ADR 0257). A copy has no builds of
// its own; it is kept by the builds that used it.
KeptStoodOn = "stood-on"
) )
// The states an artifact the mesh recorded making is in (novox/hq ADR 0251). // The states an artifact the mesh recorded making is in (novox/hq ADR 0251).
@@ -77,7 +81,7 @@ func (i *Inventory) Artifacts(ctx context.Context) ([]ArtifactState, error) {
if err != nil { if err != nil {
return nil, err return nil, err
} }
all, err := i.everyReferenceMade(ctx) all, err := i.everyReferenceRecorded(ctx)
if err != nil { if err != nil {
return nil, err return nil, err
} }
@@ -156,8 +160,9 @@ func (i *Inventory) KeepSet(ctx context.Context) (map[string][]string, error) {
// join every module ever built kept five builds' artifacts for ever — and once the store's // 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. // collector runs for real, what the keep set says is what the disk holds.
recent, err := i.store.Pool().Query(ctx, recent, err := i.store.Pool().Query(ctx,
`select made from ( `select made, built_against from (
select b.made, row_number() over (partition by b.module order by b.at desc, b.id desc) as back select b.made, b.built_against,
row_number() over (partition by b.module order by b.at desc, b.id desc) as back
from build b from build b
join module m on m.name = b.module join module m on m.name = b.module
where b.failed = '' and b.module is not null and b.module <> '' where b.failed = '' and b.module is not null and b.module <> ''
@@ -166,18 +171,26 @@ func (i *Inventory) KeepSet(ctx context.Context) (map[string][]string, error) {
return nil, err return nil, err
} }
defer recent.Close() defer recent.Close()
var stoodOn []string
for recent.Next() { for recent.Next() {
var raw []byte var raw, against []byte
if err := recent.Scan(&raw); err != nil { if err := recent.Scan(&raw, &against); err != nil {
return nil, err return nil, err
} }
var made []Artifact var made []Artifact
if err := json.Unmarshal(raw, &made); err != nil { if err := json.Unmarshal(raw, &made); err == nil {
continue for _, a := range made {
if reference := asRecorded(a.Reference); reference != "" {
because(reference, KeptByRecentBuild)
}
}
} }
for _, a := range made { var bases []string
if reference := asRecorded(a.Reference); reference != "" { if err := json.Unmarshal(against, &bases); err == nil {
because(reference, KeptByRecentBuild) for _, b := range bases {
if reference := asRecorded(b); reference != "" {
stoodOn = append(stoodOn, reference)
}
} }
} }
} }
@@ -185,6 +198,27 @@ func (i *Inventory) KeepSet(ctx context.Context) (map[string][]string, error) {
return nil, err return nil, err
} }
// **A copied base is kept by the kept builds that stood on it, and by nothing else** (novox/hq
// ADR 0257). It has no builds of its own to count five of, and no machine is ever handed a base,
// so a definition is never its reason: a definition names the image upstream, by the digest every
// copy of it shares, and read as a reason it would keep every module's copy for as long as any
// module stood on that image. A copy already let go of stays let go of.
mirrored, err := i.mirroredSet(ctx)
if err != nil {
return nil, err
}
if len(mirrored) > 0 {
collected, err := i.alreadyCollected(ctx)
if err != nil {
return nil, err
}
for _, reference := range stoodOn {
if mirrored[reference] && !collected[reference] {
because(reference, KeptStoodOn)
}
}
}
// And anything a manifest mentions — every reference ever made, the recent ones included, so a // And anything a manifest mentions — every reference ever made, the recent ones included, so a
// kept artifact carries each of its reasons and the operator reads why, not only that (novox/hq // kept artifact carries each of its reasons and the operator reads why, not only that (novox/hq
// ADR 0251). A manifest holds a reference composed with the store's address or kept bare, so the // ADR 0251). A manifest holds a reference composed with the store's address or kept bare, so the
@@ -195,6 +229,9 @@ func (i *Inventory) KeepSet(ctx context.Context) (map[string][]string, error) {
return nil, err return nil, err
} }
for _, reference := range all { for _, reference := range all {
if mirrored[reference] {
continue
}
digest := digestIn(reference) digest := digestIn(reference)
if digest == "" { if digest == "" {
// Not something the store holds by digest; nothing here can speak for it, so it // Not something the store holds by digest; nothing here can speak for it, so it
@@ -237,6 +274,72 @@ func (i *Inventory) KeptArchives(ctx context.Context) ([]string, error) {
return out, nil return out, nil
} }
// KeptImages is every image the mesh keeps, in a stated order: the manifests whose platform manifests
// a sweep spares when it lets an index of the same repository go (novox/hq ADR 0257).
func (i *Inventory) KeptImages(ctx context.Context) ([]string, error) {
keep, err := i.KeepSet(ctx)
if err != nil {
return nil, err
}
var out []string
for reference := range keep {
if path, ours := catalogue.InArtifactStore(reference); ours && strings.Contains(path, "@sha256:") {
out = append(out, reference)
}
}
sort.Strings(out)
return out, nil
}
// everyReferenceRecorded is every artifact the mesh recorded putting in its store: what each
// successful build made, and every base a build — or a person — recorded copying in (novox/hq ADR
// 0257). Oldest first, a copy placed at the moment it was recorded.
func (i *Inventory) everyReferenceRecorded(ctx context.Context) ([]string, error) {
made, err := i.everyReferenceMade(ctx)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx, `select reference from artifact_mirrored order by at asc, reference asc`)
if err != nil {
return nil, err
}
defer rows.Close()
seen := map[string]bool{}
for _, r := range made {
seen[r] = true
}
out := made
for rows.Next() {
var reference string
if err := rows.Scan(&reference); err != nil {
return nil, err
}
if reference = asRecorded(reference); reference != "" && !seen[reference] {
seen[reference] = true
out = append(out, reference)
}
}
return out, rows.Err()
}
// mirroredSet is every base the records say the mesh copied in, as recorded.
func (i *Inventory) mirroredSet(ctx context.Context) (map[string]bool, error) {
rows, err := i.store.Pool().Query(ctx, `select reference from artifact_mirrored`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]bool{}
for rows.Next() {
var reference string
if err := rows.Scan(&reference); err != nil {
return nil, err
}
out[asRecorded(reference)] = true
}
return out, rows.Err()
}
// everyReferenceMade is every artifact reference any successful build recorded, oldest build first. // everyReferenceMade is every artifact reference any successful build recorded, oldest build first.
func (i *Inventory) everyReferenceMade(ctx context.Context) ([]string, error) { func (i *Inventory) everyReferenceMade(ctx context.Context) ([]string, error) {
rows, err := i.store.Pool().Query(ctx, rows, err := i.store.Pool().Query(ctx,
@@ -331,3 +434,8 @@ func (i *Inventory) MarkCollected(ctx context.Context, references []string) erro
} }
return nil return nil
} }
// Mirrored is every base the records say the mesh copied into its store (novox/hq ADR 0257).
func (i *Inventory) Mirrored(ctx context.Context) (map[string]bool, error) {
return i.mirroredSet(ctx)
}
+5
View File
@@ -115,6 +115,11 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin
} }
for _, ref := range against[name] { for _, ref := range against[name] {
if rest, ok := strings.CutPrefix(ref, catalogue.ArtifactStoreScheme); ok { if rest, ok := strings.CutPrefix(ref, catalogue.ArtifactStoreScheme); ok {
// A copy of an upstream image is no module's artifact (novox/hq ADR 0257): its
// repository is named for the image, and its first segment names no module.
if strings.HasPrefix(rest, builder.MirrorPrefix) {
continue
}
if base, _, found := strings.Cut(rest, "/"); found { if base, _, found := strings.Cut(rest, "/"); found {
add(name, base, EdgeStandsOn) add(name, base, EdgeStandsOn)
} }
+17
View File
@@ -154,3 +154,20 @@ func TestNothingStandsOnTheBuildAgentOnlyIsBuiltByIt(t *testing.T) {
t.Errorf("everything source-built is built by the agent: %d edges", builtBy) t.Errorf("everything source-built is built by the agent: %d edges", builtBy)
} }
} }
// A copy of an upstream image is no module's artifact (novox/hq ADR 0257), even when a module happens to
// carry the name its repository starts with.
func TestACopyOfAnUpstreamImageIsNoEdge(t *testing.T) {
entries := []Entry{
{Manifest: catalogue.Manifest{Module: "route-proxy"}},
{Manifest: catalogue.Manifest{Module: "upstream"}},
}
against := map[string][]string{
"route-proxy": {catalogue.ArtifactStoreScheme + "upstream/docker.io/library/golang@sha256:a"},
}
for _, e := range dependenciesOf(entries, against, nil) {
if e.From == "route-proxy" && e.To == "upstream" {
t.Fatalf("a copy of an upstream image was read as a module's artifact: %+v", e)
}
}
}
@@ -0,0 +1,30 @@
-- A base a build copies into the artifact store is recorded where it is copied (novox/hq ADR 0257).
--
-- A build that stands on an image published elsewhere copies it into the store first (ADR 0096, 0097),
-- and until now nothing recorded the copy as something the mesh made: the build named it only in what
-- it stood on, and a build that failed after copying named it nowhere. So on 2026-10-08 the store held
-- 374 manifests of copied bases, in 27 repositories, that the sweep could never let go of, because no
-- record said the mesh had put them there (ADR 0189 §3).
--
-- One row per copy, by the reference as the mesh records it. `build` is the build that first said it
-- copied it, whether that build worked or not; `why` is a person's words when the copy was recorded by
-- hand (the `mirrors` verb), and empty otherwise. Not a foreign key to build: a build's row may be
-- written after its copy is, and a copy recorded by hand has no build.
create table artifact_mirrored (
reference text primary key,
build text,
at timestamptz not null default now(),
why text not null default ''
);
-- The copies the records already name. A build that worked kept what it stood on, and every reference
-- into a module's own `on-<argument>` repository there is a copy that build's mirror made: only the
-- mirror ever wrote a repository so named (ADR 0097). The earliest build to name each is its maker.
insert into artifact_mirrored (reference, build, at)
select distinct on (stood_on) stood_on, b.id, b.at
from build b,
jsonb_array_elements_text(case when jsonb_typeof(b.built_against) = 'array'
then b.built_against else '[]'::jsonb end) as stood_on
where stood_on ~ '^artifact-store://[a-z0-9][a-z0-9._-]*/on-[a-z0-9_]+@sha256:[0-9a-f]{64}$'
order by stood_on, b.at asc, b.id asc
on conflict (reference) do nothing;
+211
View File
@@ -0,0 +1,211 @@
package inventory
import (
"context"
"fmt"
"os"
"slices"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// A base a build copies into the store is recorded where it is copied, and kept while a kept build stood
// on it (novox/hq ADR 0257).
func mirrorRef(repository string, n int) string {
return fmt.Sprintf("%s%s@sha256:%064x", catalogue.ArtifactStoreScheme, repository, n)
}
// builtOn records a successful build of a module that made one image and stood on these bases.
func builtOn(t *testing.T, inv *Inventory, id, module string, n int, mirrored []string, against ...string) {
t.Helper()
b := aBuild(id, module, "")
b.Made = []Artifact{{Name: "app", Kind: "image", Reference: ref(module, "app", n)}}
b.Against = against
b.Mirrored = mirrored
if err := inv.RecordBuild(context.Background(), b); err != nil {
t.Fatal(err)
}
}
func stateOf(t *testing.T, inv *Inventory, reference string) ArtifactState {
t.Helper()
states, err := inv.Artifacts(context.Background())
if err != nil {
t.Fatal(err)
}
for _, s := range states {
if s.Reference == reference {
return s
}
}
return ArtifactState{Reference: reference, State: "unrecorded"}
}
func TestACopyAFailedBuildMadeIsRecordedAndNothingKeepsIt(t *testing.T) {
inv := fresh(t)
golang := mirrorRef("upstream/docker.io/library/golang", 1)
failed := aBuild("f1", "", "the recipe fetched an undeclared image")
failed.Mirrored = []string{golang}
if err := inv.RecordBuild(context.Background(), failed); err != nil {
t.Fatal(err)
}
if s := stateOf(t, inv, golang); s.State != ArtifactEligible {
t.Fatalf("a copy only a failed build made reads %+v; want eligible", s)
}
}
func TestACopyIsKeptByTheKeptBuildsThatStoodOnItAndByNothingElse(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
holding(t, inv, "proxy")
old := mirrorRef("upstream/docker.io/library/golang", 1)
current := mirrorRef("upstream/docker.io/library/golang", 2)
// The first build stood on the old copy; six more on the current one. The old copy falls out of
// the five kept builds; the current one is stood on by all five.
builtOn(t, inv, "b01", "proxy", 1, []string{old}, old)
for i := 2; i <= 7; i++ {
builtOn(t, inv, fmt.Sprintf("b%02d", i), "proxy", i, []string{current}, current)
}
if s := stateOf(t, inv, old); s.State != ArtifactEligible {
t.Fatalf("a copy no kept build stood on reads %+v", s)
}
if s := stateOf(t, inv, current); s.State != ArtifactKept || !slices.Equal(s.Why, []string{KeptStoodOn}) {
t.Fatalf("a copy the kept builds stood on reads %+v", s)
}
// A definition naming the image upstream, by the digest every copy shares, keeps no copy.
m := catalogue.Manifest{Module: "proxy", Version: "2",
Build: &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "GO_BASE", Image: "golang@" + strings.SplitN(old, "@", 2)[1]}}}}
if err := inv.RegisterModule(ctx, m, Source{Repository: "https://forge.invalid/proxy.git"}); err != nil {
t.Fatal(err)
}
if s := stateOf(t, inv, old); s.State != ArtifactEligible {
t.Fatalf("a definition naming the upstream image kept the old copy: %+v", s)
}
go_, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
if !slices.Contains(go_, old) || slices.Contains(go_, current) {
t.Fatalf("offered %v to collect", go_)
}
// Once let go of, a copy stays let go of, even should a kept build still name it.
if err := inv.MarkCollected(ctx, []string{current}); err != nil {
t.Fatal(err)
}
if s := stateOf(t, inv, current); s.State != ArtifactCollected {
t.Fatalf("a collected copy reads %+v", s)
}
}
func TestAModulesBaseIsNotKeptAsACopy(t *testing.T) {
// A base that is another module's artifact is kept by that module's own builds, as before: the
// stood-on reason is a copy's alone.
inv := fresh(t)
holding(t, inv, "app")
base := ref("go-base", "build", 1)
b := aBuild("g1", "go-base", "")
b.Made = []Artifact{{Name: "build", Kind: "image", Reference: base}}
if err := inv.RecordBuild(context.Background(), b); err != nil {
t.Fatal(err)
}
builtOn(t, inv, "a1", "app", 1, nil, base)
if s := stateOf(t, inv, base); s.State != ArtifactEligible {
t.Fatalf("a forgotten module's artifact read %+v because a build stood on it", s)
}
}
func TestAPersonRecordsACopyAndTheFirstRecordStands(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
copyRef := mirrorRef("route-proxy/on-go_base", 3)
if err := inv.RecordMirrored(ctx, "", "copied before copies were recorded", []string{copyRef}); err != nil {
t.Fatal(err)
}
if err := inv.RecordMirrored(ctx, "b9", "", []string{copyRef}); err != nil {
t.Fatal(err)
}
known, err := inv.Mirrored(ctx)
if err != nil || !known[copyRef] {
t.Fatalf("recorded copies %v %v", known, err)
}
var why string
if err := inv.store.Pool().QueryRow(ctx, `select why from artifact_mirrored where reference = $1`, copyRef).Scan(&why); err != nil {
t.Fatal(err)
}
if why != "copied before copies were recorded" {
t.Fatalf("the first record did not stand: %q", why)
}
if s := stateOf(t, inv, copyRef); s.State != ArtifactEligible {
t.Fatalf("a copy a person recorded reads %+v", s)
}
}
// The migration records the copies the build records already name: only a module's own on-<argument>
// repository, and only once, by the earliest build naming it.
func TestTheCopiesTheRecordsAlreadyNameAreRecorded(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
copyRef := mirrorRef("route-proxy/on-go_base", 1)
notACopy := ref("go-base", "build", 2)
builtOn(t, inv, "b1", "route-proxy", 1, nil, copyRef, notACopy)
builtOn(t, inv, "b2", "route-proxy", 2, nil, copyRef)
// The rows a mesh had before the table: written by builds that did not say what they copied.
if _, err := inv.store.Pool().Exec(ctx, `truncate artifact_mirrored`); err != nil {
t.Fatal(err)
}
sql, err := os.ReadFile("migrations/0081-a-mirrored-base-is-recorded.sql")
if err != nil {
t.Fatal(err)
}
_, backfill, _ := strings.Cut(string(sql), "insert into artifact_mirrored")
if _, err := inv.store.Pool().Exec(ctx, "insert into artifact_mirrored"+backfill); err != nil {
t.Fatal(err)
}
var reference, build string
var n int
if err := inv.store.Pool().QueryRow(ctx, `select count(*) from artifact_mirrored`).Scan(&n); err != nil || n != 1 {
t.Fatalf("%d copies recorded (%v); want the one", n, err)
}
if err := inv.store.Pool().QueryRow(ctx, `select reference, build from artifact_mirrored`).Scan(&reference, &build); err != nil {
t.Fatal(err)
}
if reference != copyRef || build != "b1" {
t.Fatalf("recorded %s by %s", reference, build)
}
}
// A copy let go of and copied again by a later build is held again: its record is back, and the sweep
// decides about it as before.
func TestACopyLetGoOfAndCopiedAgainIsHeldAgain(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
holding(t, inv, "proxy")
golang := mirrorRef("upstream/docker.io/library/golang", 1)
builtOn(t, inv, "b01", "proxy", 1, []string{golang}, golang)
if err := inv.MarkCollected(ctx, []string{golang}); err != nil {
t.Fatal(err)
}
if s := stateOf(t, inv, golang); s.State != ArtifactCollected {
t.Fatalf("a let-go copy reads %+v", s)
}
builtOn(t, inv, "b02", "proxy", 2, []string{golang}, golang)
if s := stateOf(t, inv, golang); s.State != ArtifactKept || !slices.Equal(s.Why, []string{KeptStoodOn}) {
t.Fatalf("a copy copied again reads %+v; want kept again", s)
}
// A person's record never undoes a collection: only a build that copied it says it is there.
if err := inv.MarkCollected(ctx, []string{golang}); err != nil {
t.Fatal(err)
}
if err := inv.RecordMirrored(ctx, "", "again", []string{golang}); err != nil {
t.Fatal(err)
}
if s := stateOf(t, inv, golang); s.State != ArtifactCollected {
t.Fatalf("a person's record brought a collected copy back: %+v", s)
}
}
+4
View File
@@ -215,6 +215,10 @@ type BuildResult struct {
// (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care. // (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care.
Against []string `json:"against,omitempty"` Against []string `json:"against,omitempty"`
// Mirrored is every base this build copied into the artifact store, as pinned (novox/hq ADR
// 0257). Carried on a failed build too: the copy is in the store whether the build worked or not.
Mirrored []string `json:"mirrored,omitempty"`
// Read is every repository this build read source from besides the module's own. A module whose // Read is every repository this build read source from besides the module's own. A module whose
// recipe packages source that lives elsewhere is affected when that repository moves, and the // recipe packages source that lives elsewhere is affected when that repository moves, and the
// manifest the mesh keeps says nothing about it (novox/hq 04-ISSUES/131). // manifest the mesh keeps says nothing about it (novox/hq 04-ISSUES/131).
+2 -1
View File
@@ -75,7 +75,8 @@
"build", "build",
"artifacts", "artifacts",
"collect", "collect",
"images" "images",
"mirrors"
], ],
"resources": [ "resources": [
{ {