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:
@@ -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
|
||||
|
||||
@@ -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})
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 " +
|
||||
|
||||
@@ -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":
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
@@ -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"},
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
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)
|
||||
store.Spare = store.SpareKeptIndexes([]string{kept})
|
||||
|
||||
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
|
||||
store.Spare = store.SpareKeptIndexes(nil)
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -9,8 +9,10 @@ package artifacts
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
@@ -24,6 +26,10 @@ type Store struct {
|
||||
Address string
|
||||
// HTTP is the client used; nil is a client with a modest timeout.
|
||||
HTTP *http.Client
|
||||
// Spare says, for one repository, which manifests a kept index there names: the platform
|
||||
// manifests an index being let go of must leave in place (novox/hq ADR 0257). Nil spares every
|
||||
// platform manifest — an index is then let go of alone, as it was before.
|
||||
Spare func(ctx context.Context, repository string) (map[string]bool, error)
|
||||
}
|
||||
|
||||
// Gone is the answer when the store does not hold it: the outcome wanted, already true.
|
||||
@@ -77,9 +83,132 @@ func (s Store) LetGo(ctx context.Context, reference string) error {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if kind == "manifests" && s.Spare != nil {
|
||||
// **An index's platform manifests go before the index** (novox/hq ADR 0257). Each was put
|
||||
// under its own digest, so each is a manifest of the repository in its own right, and the
|
||||
// store's collector keeps every manifest a repository holds: an index let go of alone frees
|
||||
// no byte of the images it names. Before, not after, so a sweep cut short leaves the index
|
||||
// to be asked about again, and with it the list of what is still to go.
|
||||
if err := s.letGoOfPlatforms(ctx, repository, digest, reference); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return s.remove(ctx, s.url(repository, kind, digest), reference)
|
||||
}
|
||||
|
||||
// letGoOfPlatforms deletes the manifests an index names, except those a kept index of the same
|
||||
// repository also names. A manifest that is not an index names none.
|
||||
func (s Store) letGoOfPlatforms(ctx context.Context, repository, digest, reference string) error {
|
||||
children, err := s.Platforms(ctx, repository, digest)
|
||||
if err != nil || len(children) == 0 {
|
||||
// Gone: there is no index to read, and the delete that follows says so.
|
||||
if errors.Is(err, Gone) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
spared, err := s.Spare(ctx, repository)
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot say which platform manifests of %s a kept index names, so %s was kept: %w",
|
||||
repository, reference, err)
|
||||
}
|
||||
for _, child := range children {
|
||||
if spared[child] {
|
||||
continue
|
||||
}
|
||||
if err := s.remove(ctx, s.url(repository, "manifests", child), reference+" (its "+child+")"); err != nil && !errors.Is(err, Gone) {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// anyManifest is every manifest the store may hold, indexes included, so it answers with the
|
||||
// document as it is.
|
||||
var anyManifest = []string{
|
||||
"application/vnd.oci.image.index.v1+json",
|
||||
"application/vnd.docker.distribution.manifest.list.v2+json",
|
||||
"application/vnd.oci.image.manifest.v1+json",
|
||||
"application/vnd.docker.distribution.manifest.v2+json",
|
||||
}
|
||||
|
||||
// Platforms is the manifests an index names, by digest; none for a manifest that is not an index.
|
||||
// Gone when the store does not hold it.
|
||||
func (s Store) Platforms(ctx context.Context, repository, digest string) ([]string, error) {
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodGet, s.url(repository, "manifests", digest), nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
request.Header.Set("Accept", strings.Join(anyManifest, ", "))
|
||||
response, err := s.client().Do(request)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer response.Body.Close()
|
||||
switch response.StatusCode {
|
||||
case http.StatusOK:
|
||||
case http.StatusNotFound:
|
||||
return nil, Gone
|
||||
default:
|
||||
return nil, fmt.Errorf("the artifact store answered %s when asked for %s@%s", response.Status, repository, digest)
|
||||
}
|
||||
var document struct {
|
||||
Manifests []struct {
|
||||
Digest string `json:"digest"`
|
||||
} `json:"manifests"`
|
||||
}
|
||||
if err := json.NewDecoder(io.LimitReader(response.Body, 4<<20)).Decode(&document); err != nil {
|
||||
return nil, fmt.Errorf("%s@%s is not a manifest: %w", repository, digest, err)
|
||||
}
|
||||
var out []string
|
||||
for _, m := range document.Manifests {
|
||||
if strings.HasPrefix(m.Digest, "sha256:") {
|
||||
out = append(out, m.Digest)
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// SpareKeptIndexes is a Spare for a sweep: the platform manifests the kept references of each
|
||||
// repository name, read from the store once per repository and remembered for the sweep. A kept
|
||||
// reference itself is spared too. A kept index the store does not hold names nothing.
|
||||
func (s Store) SpareKeptIndexes(kept []string) func(ctx context.Context, repository string) (map[string]bool, error) {
|
||||
byRepository := map[string][]string{}
|
||||
for _, reference := range kept {
|
||||
path, ours := catalogue.InArtifactStore(reference)
|
||||
if !ours {
|
||||
continue
|
||||
}
|
||||
repository, kind, digest, err := split(path)
|
||||
if err != nil || kind != "manifests" {
|
||||
continue
|
||||
}
|
||||
byRepository[repository] = append(byRepository[repository], digest)
|
||||
}
|
||||
read := map[string]map[string]bool{}
|
||||
return func(ctx context.Context, repository string) (map[string]bool, error) {
|
||||
if spared, done := read[repository]; done {
|
||||
return spared, nil
|
||||
}
|
||||
spared := map[string]bool{}
|
||||
for _, digest := range byRepository[repository] {
|
||||
spared[digest] = true
|
||||
children, err := s.Platforms(ctx, repository, digest)
|
||||
if errors.Is(err, Gone) {
|
||||
continue
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, c := range children {
|
||||
spared[c] = true
|
||||
}
|
||||
}
|
||||
read[repository] = spared
|
||||
return spared, nil
|
||||
}
|
||||
}
|
||||
|
||||
// remove asks the store to delete what is at url. Gone when it has no such thing.
|
||||
func (s Store) remove(ctx context.Context, url, what string) error {
|
||||
request, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil)
|
||||
@@ -121,3 +250,19 @@ func split(path string) (repository, kind, digest string, err error) {
|
||||
}
|
||||
return "", "", "", fmt.Errorf("%w: %q names nothing the store holds by digest", ErrNotOurs, path)
|
||||
}
|
||||
|
||||
// HoldsManifest is whether the store holds the manifest a recorded image reference names.
|
||||
func (s Store) HoldsManifest(ctx context.Context, reference string) (bool, error) {
|
||||
path, ours := catalogue.InArtifactStore(reference)
|
||||
if !ours {
|
||||
return false, fmt.Errorf("%w: %s", ErrNotOurs, reference)
|
||||
}
|
||||
repository, kind, digest, err := split(path)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if kind != "manifests" {
|
||||
return false, fmt.Errorf("%w: %s names a blob, not a manifest", ErrNotOurs, reference)
|
||||
}
|
||||
return s.has(ctx, s.url(repository, kind, digest), anyManifest...)
|
||||
}
|
||||
|
||||
@@ -77,6 +77,11 @@ type Result struct {
|
||||
// the default, and that branch is its trunk.
|
||||
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,
|
||||
// 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
|
||||
@@ -106,6 +111,19 @@ type GitCredential struct {
|
||||
func Build(ctx context.Context, run Runner, publish Publisher,
|
||||
repository, path, ref, workspace string, held map[string]string, npmrc Npmrc,
|
||||
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
|
||||
// that hand none — tests of everything but contexts — read as they did.
|
||||
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
|
||||
// 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.
|
||||
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 {
|
||||
say("bases", "copying %s into the mesh's registry", from)
|
||||
return m.MirrorImage(ctx, from, repository)
|
||||
say("bases", "copying %s into the mesh's registry as %s", 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 {
|
||||
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
|
||||
// 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,
|
||||
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 {
|
||||
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 "+
|
||||
"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 {
|
||||
return nil, nil, fmt.Errorf("%s stands on %s: %w", manifest.Module, base.Image, err)
|
||||
}
|
||||
|
||||
@@ -27,6 +27,40 @@ type Mirrorer interface {
|
||||
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 (
|
||||
mediaIndexOCI = "application/vnd.oci.image.index.v1+json"
|
||||
mediaIndexDocker = "application/vnd.docker.distribution.manifest.list.v2+json"
|
||||
@@ -86,6 +120,9 @@ func parseReference(ref string) (upstream, error) {
|
||||
type source struct {
|
||||
client *http.Client
|
||||
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
|
||||
@@ -176,6 +213,15 @@ type descriptor struct {
|
||||
// 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.
|
||||
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)
|
||||
if err != nil {
|
||||
return "", err
|
||||
@@ -195,6 +241,17 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
|
||||
}
|
||||
}
|
||||
src := &source{client: r.client()}
|
||||
if former != "" && former != repository && strings.HasPrefix(where.reference, "sha256:") {
|
||||
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+former+"/manifests/"+where.reference, manifestAccept)
|
||||
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)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("copying %s into %s/%s: %w", from, r.Address, repository, err)
|
||||
@@ -286,16 +343,14 @@ func (r Registry) copyBlob(ctx context.Context, src *source, where upstream, dig
|
||||
} else if there {
|
||||
return nil
|
||||
}
|
||||
response, err := src.get(ctx, where.base+"/v2/"+where.repository+"/blobs/"+digest, "")
|
||||
if err != nil {
|
||||
return err
|
||||
// **Mounted where this registry already holds it** (novox/hq ADR 0257): a blob is stored once
|
||||
// whichever repositories link it, so a mount moves no bytes. A registry that cannot mount answers
|
||||
// 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()
|
||||
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)
|
||||
start, err := http.NewRequestWithContext(ctx, http.MethodPost, uploads, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -304,9 +359,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)
|
||||
}
|
||||
begun.Body.Close()
|
||||
if begun.StatusCode == http.StatusCreated && src.mountFrom != "" {
|
||||
return nil
|
||||
}
|
||||
if begun.StatusCode != http.StatusAccepted {
|
||||
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")
|
||||
if location == "" {
|
||||
return fmt.Errorf("%s accepted an upload and said nowhere to put it", base)
|
||||
|
||||
@@ -0,0 +1,245 @@
|
||||
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
|
||||
}
|
||||
|
||||
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 == "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.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)
|
||||
}
|
||||
}
|
||||
@@ -72,7 +72,7 @@ func TestAnIncompleteBaseIsRefused(t *testing.T) {
|
||||
}
|
||||
|
||||
// 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")
|
||||
}
|
||||
|
||||
@@ -86,17 +86,20 @@ func TestADeclaredVendorImageIsCopiedInAndHandedToTheRecipe(t *testing.T) {
|
||||
},
|
||||
}
|
||||
var asked []string
|
||||
args, _, err := standingOn(context.Background(), manifest, nil, func(_ context.Context, from, repository string) (string, error) {
|
||||
asked = append(asked, from+" -> "+repository)
|
||||
args, _, err := standingOn(context.Background(), manifest, nil, func(_ context.Context, from, repository, former string) (string, error) {
|
||||
asked = append(asked, from+" -> "+repository+" (from "+former+")")
|
||||
return "127.0.0.1:5000/" + repository + "@sha256:" + strings.Repeat("d", 64), nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(asked) != 1 || asked[0] != "quay.io/minio/mc@sha256:"+strings.Repeat("c", 64)+" -> minio/on-mc_base" {
|
||||
t.Fatalf("the image was not copied under the module's repository: %v", asked)
|
||||
// One repository per upstream image, named for the image (novox/hq ADR 0257); the module's
|
||||
// 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)
|
||||
}
|
||||
// Unpinned, it is refused: a tag is what somebody else can move.
|
||||
|
||||
@@ -457,6 +457,19 @@ var ControllerVerbs = []Verb{
|
||||
"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)",
|
||||
}, 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. Records only: nothing is deleted " +
|
||||
"(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 " +
|
||||
"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 " +
|
||||
|
||||
@@ -39,6 +39,9 @@ type Build struct {
|
||||
// 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).
|
||||
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
|
||||
// 04-ISSUES/131), at the ref it read.
|
||||
Read []ReadRepository
|
||||
@@ -120,7 +123,30 @@ func (i *Inventory) RecordBuild(ctx context.Context, b Build) error {
|
||||
on conflict (id) do nothing`,
|
||||
b.ID, b.Repository, b.Ref, module, b.Commit, b.On, b.Failed, made,
|
||||
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.
|
||||
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
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Builds is what has happened lately, newest first.
|
||||
|
||||
@@ -37,6 +37,10 @@ const (
|
||||
// KeptUnaddressable: it names no digest, so nothing here can speak for it, and it is kept rather
|
||||
// than guessed about.
|
||||
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).
|
||||
@@ -77,7 +81,7 @@ func (i *Inventory) Artifacts(ctx context.Context) ([]ArtifactState, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
all, err := i.everyReferenceMade(ctx)
|
||||
all, err := i.everyReferenceRecorded(ctx)
|
||||
if err != nil {
|
||||
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
|
||||
// collector runs for real, what the keep set says is what the disk holds.
|
||||
recent, err := i.store.Pool().Query(ctx,
|
||||
`select made from (
|
||||
select b.made, row_number() over (partition by b.module order by b.at desc, b.id desc) as back
|
||||
`select made, built_against from (
|
||||
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
|
||||
join module m on m.name = 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
|
||||
}
|
||||
defer recent.Close()
|
||||
var stoodOn []string
|
||||
for recent.Next() {
|
||||
var raw []byte
|
||||
if err := recent.Scan(&raw); err != nil {
|
||||
var raw, against []byte
|
||||
if err := recent.Scan(&raw, &against); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var made []Artifact
|
||||
if err := json.Unmarshal(raw, &made); err != nil {
|
||||
continue
|
||||
if err := json.Unmarshal(raw, &made); err == nil {
|
||||
for _, a := range made {
|
||||
if reference := asRecorded(a.Reference); reference != "" {
|
||||
because(reference, KeptByRecentBuild)
|
||||
}
|
||||
}
|
||||
}
|
||||
for _, a := range made {
|
||||
if reference := asRecorded(a.Reference); reference != "" {
|
||||
because(reference, KeptByRecentBuild)
|
||||
var bases []string
|
||||
if err := json.Unmarshal(against, &bases); err == nil {
|
||||
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
|
||||
}
|
||||
|
||||
// **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
|
||||
// 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
|
||||
@@ -195,6 +229,9 @@ func (i *Inventory) KeepSet(ctx context.Context) (map[string][]string, error) {
|
||||
return nil, err
|
||||
}
|
||||
for _, reference := range all {
|
||||
if mirrored[reference] {
|
||||
continue
|
||||
}
|
||||
digest := digestIn(reference)
|
||||
if digest == "" {
|
||||
// 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
|
||||
}
|
||||
|
||||
// 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.
|
||||
func (i *Inventory) everyReferenceMade(ctx context.Context) ([]string, error) {
|
||||
rows, err := i.store.Pool().Query(ctx,
|
||||
@@ -331,3 +434,8 @@ func (i *Inventory) MarkCollected(ctx context.Context, references []string) erro
|
||||
}
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -115,6 +115,11 @@ func dependenciesOf(entries []Entry, against map[string][]string, read map[strin
|
||||
}
|
||||
for _, ref := range against[name] {
|
||||
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 {
|
||||
add(name, base, EdgeStandsOn)
|
||||
}
|
||||
|
||||
@@ -154,3 +154,20 @@ func TestNothingStandsOnTheBuildAgentOnlyIsBuiltByIt(t *testing.T) {
|
||||
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;
|
||||
@@ -0,0 +1,181 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -215,6 +215,10 @@ type BuildResult struct {
|
||||
// (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care.
|
||||
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
|
||||
// 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).
|
||||
|
||||
+2
-1
@@ -75,7 +75,8 @@
|
||||
"build",
|
||||
"artifacts",
|
||||
"collect",
|
||||
"images"
|
||||
"images",
|
||||
"mirrors"
|
||||
],
|
||||
"resources": [
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user