Record the bases a build copies, keep them by the builds that stood on them, and copy each image once

A copied base was named only in what a build stood on, and nowhere when the build failed,
so the store's sweep could never let one go (hq issue 321). One repository per upstream
image stops each module asking the public registry for the same image again, and letting
an index go now takes its own platform manifests, which otherwise kept every byte. A
person can record the copies no record names through the new mirrors verb (hq ADR 0257).

The forge test fix is the same commit as on feat/plain-notifications: main fails without it.
This commit is contained in:
jochen
2026-10-08 15:32:10 +02:00
parent 34611c8e38
commit 7edfb56e7a
25 changed files with 1472 additions and 33 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,
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
// failed on its own — in the moment the kill arrived says what it did, and the kill is refused.
killed := false
+3
View File
@@ -160,6 +160,9 @@ func buildFrom(result link.BuildResult) inventory.Build {
for _, ref := range result.Against {
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 {
kept.Read = append(kept.Read, inventory.ReadRepository{Repository: r.Repository, Ref: r.Ref})
}
+14
View File
@@ -76,6 +76,20 @@ 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
// 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.
// **An index goes with its platform manifests, and a kept index keeps its own** (novox/hq ADR
// 0257): which of a repository's platform manifests a kept index names is read from the store
// before anything there is let go of. A keep set that cannot be read stops the sweep before it
// deletes anything, as a kept archive that cannot be held does.
if store.Spare == nil && inv != nil {
images, err := inv.KeptImages(ctx)
if err != nil {
r.Stopped = fmt.Sprintf("the images the mesh keeps could not be read, so nothing was let go: %v", err)
r.Left = len(references)
return r
}
store.Spare = store.SpareKeptIndexes(images)
}
wrote, missing, err := holdKept(within, store, kept)
r.HoldersWritten, r.Missing = wrote, missing
if err != nil {
+4
View File
@@ -86,6 +86,10 @@ var handActVerbs = []handActVerb{
// 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 " +
"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 " +
"person starts (ADR 0236)"},
{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:])
case "images":
return imagesCommand(ctx, args[1:])
case "mirrors":
return mirrorsCommand(ctx, args[1:])
case "plans":
return plansCommand(ctx, args[1:])
case "delivery":
+221
View File
@@ -0,0 +1,221 @@
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"`
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 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
}
if *confirm {
f.record(ctx, "mirrors", []string{"--record", *record})
}
answer, err = recordMirrors(ctx, inv, artifacts.Store{Address: address}, splitReferences(*record), known,
strings.TrimSpace(*f.why), *confirm)
if err != nil {
return err
}
}
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)")
}
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
}
+138
View File
@@ -0,0 +1,138 @@
package main
import (
"fmt"
"net/http"
"net/http/httptest"
"slices"
"strings"
"testing"
"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 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)
}
}
}
@@ -133,6 +133,10 @@ func (f *fakeStore) serve(t *testing.T) artifacts.Store {
w.WriteHeader(http.StatusOK)
case r.Method == http.MethodHead:
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:
w.Header().Set("Location", "/upload/x?state=1")
w.WriteHeader(http.StatusAccepted)
+27 -1
View File
@@ -605,6 +605,30 @@ func (a *verbArguments) commandLine() ([]string, error) {
return nil, err
}
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":
argv := []string{"data", "--json"}
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,
"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).
"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.
"delivery": true}
@@ -835,6 +859,8 @@ func repairingCommand(argv []string) string {
return "cleanup delete"
case argv[0] == "collect" && slices.Contains(argv, "--confirm"):
return "collect"
case argv[0] == "mirrors" && slices.Contains(argv, "--confirm"):
return "mirrors"
}
return ""
}
@@ -284,6 +284,7 @@ var accountedFlags = map[string]map[string]string{
"artifacts": {"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"},
"mirrors": {"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"},
"healers": {"json": "set by the verb: the answer is data"},