Merge pull request 'Three verbs for the registries: artifacts, collect, images (hq ADR 0251)' (#140) from feat/registry-verbs into main

This commit was merged in pull request #140.
This commit is contained in:
2026-10-08 11:05:35 +00:00
12 changed files with 1116 additions and 127 deletions
+127 -76
View File
@@ -17,14 +17,126 @@ import (
// moment the keep set moved, so it is the moment to say what may go — and it needs no timer of
// its own. Reclaiming the bytes is the store's own nightly step; this only decides.
//
// **And run when a person asks** (novox/hq ADR 0251): the `collect` verb is the same sweep, with
// bounds of its own. One implementation, so what a person runs is exactly what a build runs.
//
// Never fatal to a build. The build succeeded, the module is registered, and a store that could
// not be reached is a thing to say rather than a reason to undo any of that. The next build asks
// again, and the references it could not collect are still uncollected, so nothing is lost by
// having failed.
// sweepBounds is how much one sweep may ask of the store.
type sweepBounds struct {
// most is how many artifacts it lets go of at most.
most int
// budget is the longest it keeps its caller waiting.
budget time.Duration
}
// afterBuild are the bounds of the sweep inside somebody's build.
//
// **Bounded, because this runs inside somebody's build.** The first sweep of a mesh that has never
// collected has the whole history to get through, and a person waiting on `build` should not pay for
// it. What is left over is simply offered again next time — builds are frequent, and the point is
// that the store stops growing, not that it empties tonight.
var afterBuild = sweepBounds{most: mostPerSweep, budget: sweepBudget}
// sweepResult is what one sweep did.
type sweepResult struct {
// HoldersWritten is how many kept archives the store now holds by a manifest it did not before.
HoldersWritten int
// Missing is how many kept archives the store does not have at all.
Missing int
// Held is whether every kept archive could be asked to be held; false stops the sweep before it
// lets anything go.
Held bool
// LetGo is what the store let go of, and the record now says collected.
LetGo []string
// Skipped is how many references the sweep will not address (novox/hq issue 226).
Skipped int
// Left is how many eligible references were not asked about this time.
Left int
// Bounded is whether it stopped at its bounds rather than at a refusal.
Bounded bool
// Stopped says why the sweep stopped before the end; empty when it did not.
Stopped string
}
// sweep holds every kept archive, then asks the store to let go of each reference in order, records
// what it let go of, and says what it did.
func sweep(ctx context.Context, inv *inventory.Inventory, store artifacts.Store, references, kept []string,
bounds sweepBounds) sweepResult {
within, stop := context.WithTimeout(ctx, bounds.budget)
defer stop()
var r sweepResult
// **Hold before letting go** (novox/hq issue 253). The store's collector keeps only what a
// manifest names, and archives were published as bare blobs, so every kept archive is first
// held by its manifest — which backfills the ones published before holders, a few at a time
// 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.
wrote, missing, err := holdKept(within, store, kept)
r.HoldersWritten, r.Missing = wrote, missing
if err != nil {
r.Stopped = fmt.Sprintf("not every kept archive could be held, so nothing was let go: %v", err)
r.Left = len(references)
return r
}
r.Held = true
for i, reference := range references {
if i >= bounds.most || within.Err() != nil {
r.Left = len(references) - i
r.Bounded = true
if i >= bounds.most {
r.Stopped = fmt.Sprintf("asked about %d, the most one sweep asks", bounds.most)
} else {
r.Stopped = fmt.Sprintf("ran out of its %s", bounds.budget)
}
break
}
err := store.LetGo(within, reference)
if err == nil || errors.Is(err, artifacts.Gone) {
// Gone is the outcome wanted, already true. Recorded so the next sweep does not ask
// again for ever.
r.LetGo = append(r.LetGo, reference)
continue
}
if errors.Is(err, artifacts.ErrNotOurs) {
// **A fact about this record, so this record is skipped** (novox/hq issue 226). Not
// marked collected — the mesh did not remove it and should not claim to — and not a
// reason to stop, because the store was never asked. One of these at the front of
// the oldest-first order ended every sweep until this.
r.Skipped++
if r.Skipped == 1 {
fmt.Fprintf(os.Stderr, "the sweep will not address %s and went on: %v\n", reference, err)
}
continue
}
// **Stopped at the first refusal by the STORE, not pushed through.** A store that refuses
// one refuses all of them — deletion disabled, the store down, the network gone — so
// going on would be a hundred identical failures and a hundred identical log lines in
// front of whoever was building something.
r.Stopped = fmt.Sprintf("the artifact store kept %s, so nothing more was asked of it: %v", reference, err)
r.Left = len(references) - i
break
}
if len(r.LetGo) > 0 {
// Recorded outside `within`: the deletions happened, and losing the record of them because
// the sweep ran out of budget would mean asking about them again for ever.
if err := inv.MarkCollected(ctx, r.LetGo); err != nil {
r.Stopped = fmt.Sprintf("the store let go of %d artifact(s) and the record of it did not keep: %v",
len(r.LetGo), err)
}
}
return r
}
// collect asks the store to let go of everything the mesh made and no longer keeps, and records
// what it let go of. Says what it did and what it could not; returns nothing, because nothing
// upstream should branch on it.
// what it let go of — after a build. Says what it did and what it could not; returns nothing,
// because nothing upstream should branch on it.
func collect(ctx context.Context, inv *inventory.Inventory) {
references, err := inv.ToCollect(ctx)
if err != nil {
@@ -55,86 +167,25 @@ func collect(ctx context.Context, inv *inventory.Inventory) {
return
}
// **Bounded, because this runs inside somebody's build.** The first sweep of a mesh that has
// never collected has the whole history to get through, and a person waiting on `build` should
// not pay for it. Two bounds, and what is left over is simply offered again next time —
// builds are frequent, and the point is that the store stops growing, not that it empties
// tonight.
within, stop := context.WithTimeout(ctx, sweepBudget)
defer stop()
store := artifacts.Store{Address: address}
// **Hold before letting go** (novox/hq issue 253). The store's collector keeps only what a
// manifest names, and archives were published as bare blobs, so every kept archive is first
// held by its manifest — which backfills the ones published before holders, a few at a time
// 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.
wrote, missing, err := holdKept(within, store, kept)
if wrote > 0 {
fmt.Fprintf(os.Stderr, "the artifact store now holds %d more kept archive(s) by a manifest\n", wrote)
r := sweep(ctx, inv, artifacts.Store{Address: address}, references, kept, afterBuild)
if r.HoldersWritten > 0 {
fmt.Fprintf(os.Stderr, "the artifact store now holds %d more kept archive(s) by a manifest\n", r.HoldersWritten)
}
if missing > 0 {
if r.Missing > 0 {
fmt.Fprintf(os.Stderr, "%d archive(s) the mesh keeps are not in the artifact store at all; "+
"`collection` lists them\n", missing)
"`collection` lists them\n", r.Missing)
}
if err != nil {
fmt.Fprintf(os.Stderr, "not every kept archive could be held, so nothing was let go: %v\n", err)
return
if len(r.LetGo) > 0 {
fmt.Fprintf(os.Stderr, "the artifact store let go of %d artifact(s) the mesh no longer keeps\n", len(r.LetGo))
}
var done []string
var left, skipped int
for i, reference := range references {
if i >= mostPerSweep || within.Err() != nil {
left = len(references) - i
break
}
err := store.LetGo(within, reference)
if err == nil || errors.Is(err, artifacts.Gone) {
// Gone is the outcome wanted, already true. Recorded so the next sweep does not ask
// again for ever.
done = append(done, reference)
continue
}
if errors.Is(err, artifacts.ErrNotOurs) {
// **A fact about this record, so this record is skipped** (novox/hq issue 226). Not
// marked collected — the mesh did not remove it and should not claim to — and not a
// reason to stop, because the store was never asked. One of these at the front of
// the oldest-first order ended every sweep until this.
skipped++
if skipped == 1 {
fmt.Fprintf(os.Stderr,
"the sweep will not address %s and went on: %v\n", reference, err)
}
continue
}
// **Stopped at the first refusal by the STORE, not pushed through.** A store that refuses
// one refuses all of them — deletion disabled, the store down, the network gone — so
// going on would be a hundred identical failures and a hundred identical log lines in
// front of whoever was building something.
fmt.Fprintf(os.Stderr, "the artifact store kept %s, so nothing more was asked of it: %v\n",
reference, err)
left = len(references) - i
break
if r.Stopped != "" && !r.Bounded {
fmt.Fprintln(os.Stderr, r.Stopped)
}
if len(done) > 0 {
// Recorded outside `within`: the deletions happened, and losing the record of them because
// the sweep ran out of budget would mean asking about them again for ever.
if err := inv.MarkCollected(ctx, done); err != nil {
fmt.Fprintf(os.Stderr, "the store let go of %d artifact(s) and the record of it did not keep: %v\n",
len(done), err)
return
}
fmt.Fprintf(os.Stderr, "the artifact store let go of %d artifact(s) the mesh no longer keeps\n",
len(done))
if r.Held && r.Left > 0 {
fmt.Fprintf(os.Stderr, "%d more to collect; the next build asks again\n", r.Left)
}
if left > 0 {
fmt.Fprintf(os.Stderr, "%d more to collect; the next build asks again\n", left)
}
if skipped > 0 {
fmt.Fprintf(os.Stderr, "%d artifact(s) the sweep will not address were skipped\n", skipped)
if r.Skipped > 0 {
fmt.Fprintf(os.Stderr, "%d artifact(s) the sweep will not address were skipped\n", r.Skipped)
}
}
+4
View File
@@ -79,6 +79,10 @@ var handActVerbs = []handActVerb{
{Verb: "retire approve", Decision: "nothing is retired past its bound without a person (ADR 0230)"},
{Verb: "retire reject", Decision: "keeping a consumer active is a person's word (ADR 0230)"},
{Verb: "cleanup delete", Decision: "nothing retired is deleted without a person (ADR 0230)"},
// The sweep run on a person's word rather than after a build: the same decision the records make, at
// 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)"},
{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 " +
+11
View File
@@ -93,6 +93,13 @@ func run() error {
return pauseCommand(ctx, args[0], args[1:])
case "collection":
return collectionCommand(ctx, args[1:])
// What the records say the registries may keep, asked on demand (novox/hq ADR 0251).
case "artifacts":
return artifactsCommand(ctx, args[1:])
case "collect":
return collectCommand(ctx, args[1:])
case "images":
return imagesCommand(ctx, args[1:])
case "plans":
return plansCommand(ctx, args[1:])
case "delivery":
@@ -240,6 +247,10 @@ func usage() {
cleanup [list] [--json] every retired consumer per provider: age, size, why
cleanup delete <node> <module> <consumer> --why <text> the provider deletes that one retired consumer
cleanup delete --older-than <days> --why <text> [--confirm] list those older; delete only with --confirm
artifacts [--json] [--repository <r>] [--collected] every artifact the mesh made: kept and why, eligible, collected (ADR 0251)
collect [--json] [--most <n>] what the store's sweep would let go of, now; nothing is done
collect --confirm --why <text> [--most <n>] run the sweep now: hold every kept archive, let go of the eligible
images <node> [--json] every container image the machine's declaration names, now and as last sent
seats [--json] every seat this mesh defines, what it delivers, and who holds it
seat rename <from> <to> rename a seat; its former name still resolves (ADR 0122)
seat <name> --to <node>/<module> hand a seat to that assignment as one act; never empty in between (ADR 0131)
+443
View File
@@ -0,0 +1,443 @@
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"os"
"slices"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// What the records say the registries may keep, asked on demand (novox/hq ADR 0251, to-be 51).
//
// Three commands, each a verb of the mesh-controller seat:
//
// - `artifacts` — every reference the mesh recorded making, with its state: what the store's own
// tools set beside what the store holds;
// - `collect` — the sweep that runs after a build, run when a person asks, a dry run without
// --confirm;
// - `images` — every container image a machine's declaration names, now and as last sent: what
// a machine keeps when it prunes its images.
// artifactEntry is one recorded reference as `artifacts` answers it. Compact: the answer carries every
// kept and eligible reference the mesh has, in one bus message.
type artifactEntry struct {
Reference string `json:"reference"`
Repository string `json:"repository"`
Digest string `json:"digest"`
Kind string `json:"kind"`
State string `json:"state"`
Why []string `json:"why,omitempty"`
}
type artifactCounts struct {
Kept int `json:"kept"`
Eligible int `json:"eligible"`
Collected int `json:"collected"`
}
type artifactsAnswer struct {
Store string `json:"store"`
KeptBuilds int `json:"kept_builds"`
Counts artifactCounts `json:"counts"`
References []artifactEntry `json:"references"`
}
// artifactEntryOf reads a recorded reference into its repository, digest and kind; false for one that
// is not into the mesh's artifact store or names nothing by digest there.
func artifactEntryOf(s inventory.ArtifactState) (artifactEntry, bool) {
path, ours := catalogue.InArtifactStore(s.Reference)
if !ours {
return artifactEntry{}, false
}
e := artifactEntry{Reference: s.Reference, State: s.State, Why: s.Why}
if repository, digest, ok := strings.Cut(path, "@sha256:"); ok {
e.Repository, e.Digest, e.Kind = repository, "sha256:"+digest, "image"
return e, true
}
if repository, digest, ok := strings.Cut(path, "/blobs/sha256:"); ok {
e.Repository, e.Digest, e.Kind = repository, "sha256:"+digest, "archive"
return e, true
}
return artifactEntry{}, false
}
// artifactsOf is the answer for these states: counted whole, listed narrowed.
//
// **Collected references are listed only when asked for** (novox/hq ADR 0251, issue 314). They are
// history: every artifact the mesh ever let go stays collected for ever, so their number only grows,
// and a verb's answer is one bus message — an answer that outgrows it is lost on its way back. Kept
// and eligible are bounded by what the mesh holds now; the collected are counted, always.
func artifactsOf(states []inventory.ArtifactState, repository string, withCollected bool) artifactsAnswer {
a := artifactsAnswer{KeptBuilds: inventory.KeptBuilds, References: []artifactEntry{}}
for _, s := range states {
e, ours := artifactEntryOf(s)
if !ours {
continue
}
switch s.State {
case inventory.ArtifactKept:
a.Counts.Kept++
case inventory.ArtifactEligible:
a.Counts.Eligible++
case inventory.ArtifactCollected:
a.Counts.Collected++
}
if repository != "" && e.Repository != repository {
continue
}
if s.State == inventory.ArtifactCollected && !withCollected {
continue
}
a.References = append(a.References, e)
}
return a
}
func artifactsCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("artifacts", flag.ContinueOnError)
asJSON := set.Bool("json", false, "answer as JSON")
repository := set.String("repository", "", "only this repository's references (<module>/<artifact>); the counts stay whole")
collected := set.Bool("collected", false, "list the collected references too")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 0 {
return errors.New("artifacts [--json] [--repository <module>/<artifact>] [--collected]")
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
states, err := open.inventory.Artifacts(ctx)
if err != nil {
return err
}
answer := artifactsOf(states, strings.TrimSpace(*repository), *collected)
if shelf, err := open.inventory.Catalogue(ctx); err == nil {
answer.Store, _ = artifactStoreAddress(ctx, open.inventory, shelf, "")
}
if *asJSON {
encoder := json.NewEncoder(os.Stdout)
encoder.SetIndent("", " ")
return encoder.Encode(answer)
}
fmt.Printf("artifact store %s\nkept %d\neligible %d\ncollected %d\n",
storeOrNone(answer.Store), answer.Counts.Kept, answer.Counts.Eligible, answer.Counts.Collected)
for _, e := range answer.References {
why := ""
if len(e.Why) > 0 {
why = " (" + strings.Join(e.Why, ", ") + ")"
}
fmt.Printf(" %-9s %s%s\n", e.State, e.Reference, why)
}
return nil
}
func storeOrNone(s string) string {
if s == "" {
return "(not on the network)"
}
return s
}
// collectAnswer is what `collect` did, or — a dry run — would do.
type collectAnswer struct {
DryRun bool `json:"dry_run"`
Store string `json:"store"`
KeptArchives int `json:"kept_archives"`
// Held is how many kept archives the store holds by a manifest: asked with HEADs in a dry run,
// after the holding in a real one.
Held int `json:"held"`
Unheld int `json:"unheld"`
HoldersWritten int `json:"holders_written"`
Missing int `json:"missing"`
Eligible int `json:"eligible"`
// WouldLetGo is the eligible references, oldest first, at most mostListed of them.
WouldLetGo []string `json:"would_let_go"`
NotListed int `json:"not_listed,omitempty"`
LetGo []string `json:"let_go"`
Skipped int `json:"skipped"`
Left int `json:"left"`
Stopped string `json:"stopped,omitempty"`
}
// mostListed is how many references `collect` lists by name; the rest are counted.
const mostListed = 1000
// onDemand are the bounds of a sweep a person asks for: more than after a build, inside a module's
// sixty-second ask (novox/hq ADR 0251).
const (
collectMostDefault = 500
collectMostAtMost = 5000
collectBudget = 45 * time.Second
)
// causeCollect is the cause `collect` records when none is given.
const causeCollect = "collect-on-demand"
func collectCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("collect", flag.ContinueOnError)
asJSON := set.Bool("json", false, "answer as JSON")
confirm := set.Bool("confirm", false, "let go of what is eligible, rather than only say what would go")
most := set.Int("most", collectMostDefault, "let go of at most this many")
f := addHandActFlags(set)
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 0 {
return errors.New("collect [--json] [--confirm --why <text>] [--most <n>]")
}
if *confirm {
// A deletion a person asks for says why, before anything is asked of the store.
if err := f.require("collect"); err != nil {
return err
}
}
if *most < 1 {
return fmt.Errorf("collect: --most is at least 1, not %d", *most)
}
bounds := sweepBounds{most: min(*most, collectMostAtMost), budget: collectBudget}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
references, err := inv.ToCollect(ctx)
if err != nil {
return err
}
kept, err := inv.KeptArchives(ctx)
if err != nil {
return err
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
return err
}
address, err := artifactStoreAddress(ctx, inv, shelf, "")
if err != nil {
return err
}
if *confirm {
if strings.TrimSpace(*f.cause) == "" {
*f.cause = causeCollect
}
f.record(ctx, "collect", []string{"--most", fmt.Sprint(bounds.most)})
}
answer := runCollect(ctx, inv, artifacts.Store{Address: address}, references, kept, *confirm, bounds)
if *asJSON {
encoder := json.NewEncoder(os.Stdout)
encoder.SetIndent("", " ")
return encoder.Encode(answer)
}
printCollect(answer)
return nil
}
// runCollect is `collect` against a store: a dry run asks only HEADs, a real one is the sweep.
func runCollect(ctx context.Context, inv *inventory.Inventory, store artifacts.Store, references, kept []string,
real bool, bounds sweepBounds) collectAnswer {
a := collectAnswer{DryRun: !real, Store: store.Address, KeptArchives: len(kept), Eligible: len(references),
WouldLetGo: []string{}, LetGo: []string{}}
listed := references
if len(listed) > mostListed {
listed, a.NotListed = listed[:mostListed], len(listed)-mostListed
}
a.WouldLetGo = append(a.WouldLetGo, listed...)
if store.Address == "" {
a.Stopped = "this mesh has no artifact store on its network: nothing was asked of it"
a.Left = len(references)
return a
}
if !real {
// Reads only: each kept archive asked about with HEADs, as `collection` does.
report := collectionReport{Unheld: []string{}, Missing: []string{}}
unasked, stopped := askHeld(ctx, store, kept, &report)
a.Held, a.Unheld, a.Missing = report.Held, len(report.Unheld), len(report.Missing)
if unasked > 0 {
a.Stopped = fmt.Sprintf("%d kept archive(s) not asked about: %s", unasked, stopped)
}
return a
}
r := sweep(ctx, inv, store, references, kept, bounds)
a.HoldersWritten, a.Missing, a.Skipped, a.Left, a.Stopped = r.HoldersWritten, r.Missing, r.Skipped, r.Left, r.Stopped
if r.Held {
a.Held = len(kept) - r.Missing
}
a.LetGo = append(a.LetGo, r.LetGo...)
return a
}
func printCollect(a collectAnswer) {
if a.DryRun {
fmt.Println("a dry run: nothing was held and nothing was let go (--confirm --why <text> to collect)")
}
fmt.Printf("artifact store %s\n", storeOrNone(a.Store))
fmt.Printf("kept archives %d (held %d, unheld %d, missing %d)\n", a.KeptArchives, a.Held, a.Unheld, a.Missing)
if a.HoldersWritten > 0 {
fmt.Printf("holders written %d\n", a.HoldersWritten)
}
fmt.Printf("eligible to let go %d\n", a.Eligible)
if !a.DryRun {
fmt.Printf("let go %d (skipped %d, left %d)\n", len(a.LetGo), a.Skipped, a.Left)
}
if a.Stopped != "" {
fmt.Printf("stopped: %s\n", a.Stopped)
}
for _, reference := range a.WouldLetGo {
fmt.Printf(" %s\n", reference)
}
if a.NotListed > 0 {
fmt.Printf(" … and %d more not listed\n", a.NotListed)
}
}
// declaredImage is one image a machine's declaration names.
type declaredImage struct {
Image string `json:"image"`
Resources []string `json:"resources"`
// In says which declarations name it: "declaration" (what the mesh would send now), "sent" (what
// the machine was last sent).
In []string `json:"in"`
}
type imagesAnswer struct {
Node string `json:"node"`
Images []declaredImage `json:"images"`
// SentKnown is whether what the machine was last sent is kept with its images; false when it is
// not kept, or was kept before the images were.
SentKnown bool `json:"sent_known"`
}
// imagesOf is every container image of the declaration about to be sent and of the summary of the one
// last sent.
func imagesOf(node string, body []byte, sent []sentResource, sentKept bool) (imagesAnswer, error) {
a := imagesAnswer{Node: node, Images: []declaredImage{}, SentKnown: sentKept}
byImage := map[string]*declaredImage{}
add := func(image, resource, in string) {
if image == "" {
return
}
d := byImage[image]
if d == nil {
d = &declaredImage{Image: image}
byImage[image] = d
}
if !slices.Contains(d.Resources, resource) {
d.Resources = append(d.Resources, resource)
}
if !slices.Contains(d.In, in) {
d.In = append(d.In, in)
}
}
if body != nil {
var envelope struct {
Resources []map[string]any `json:"resources"`
}
if err := json.Unmarshal(body, &envelope); err != nil {
return a, fmt.Errorf("a declaration that is not one: %w", err)
}
for _, r := range envelope.Resources {
if r["type"] != "container" {
continue
}
image, _ := r["image"].(string)
add(image, fmt.Sprint(r["id"]), "declaration")
}
}
for _, s := range sent {
if s.Type != "container" {
continue
}
if s.Image == "" {
// Kept before the summary carried images: what was sent cannot be said.
a.SentKnown = false
continue
}
add(s.Image, s.ID, "sent")
}
for _, d := range byImage {
sort.Strings(d.Resources)
a.Images = append(a.Images, *d)
}
sort.Slice(a.Images, func(i, j int) bool { return a.Images[i].Image < a.Images[j].Image })
return a, nil
}
func imagesCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("images", flag.ContinueOnError)
asJSON := set.Bool("json", false, "answer as JSON")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) != 1 {
return errors.New("images <node> [--json]")
}
node := positionals[0]
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
if _, err := open.inventory.NodeByName(ctx, node); err != nil {
return fmt.Errorf("%s: %w", node, err)
}
// What the mesh would send now, composed as `plan` composes it.
var body []byte
plan, settings, err := planFor(ctx, open, node)
if err != nil {
return err
}
if len(plan.Modules) > 0 {
declared, err := declarationFor(ctx, open, node, plan, settings)
if err != nil {
return err
}
if body, err = declared.Body(); err != nil {
return err
}
}
// And what it was last sent, as `plan --diff` compares against.
prev, err := open.inventory.SentSummary(ctx, node)
if err != nil {
return err
}
var sent []sentResource
if len(prev) > 0 {
if err := json.Unmarshal(prev, &sent); err != nil {
return fmt.Errorf("%s: what it was last sent is kept in a form this version does not read: %w", node, err)
}
}
answer, err := imagesOf(node, body, sent, len(prev) > 0)
if err != nil {
return err
}
if *asJSON {
encoder := json.NewEncoder(os.Stdout)
encoder.SetIndent("", " ")
return encoder.Encode(answer)
}
for _, d := range answer.Images {
fmt.Printf("%s (%s; %s)\n", d.Image, strings.Join(d.In, ", "), strings.Join(d.Resources, ", "))
}
if !answer.SentKnown {
fmt.Printf("%s: which images it was last sent is not kept — the next push keeps it\n", node)
}
return nil
}
+299
View File
@@ -0,0 +1,299 @@
package main
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"slices"
"strings"
"sync"
"testing"
"time"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// What the records say the registries may keep, asked on demand (novox/hq ADR 0251).
func TestTheRegistryVerbsComposeTheirCommandsAndRefuseWhatTheyDoNotTake(t *testing.T) {
for _, c := range []struct {
verb string
args map[string]any
want string
}{
{"artifacts", map[string]any{}, "artifacts --json"},
{"artifacts", map[string]any{"repository": "web/app", "collected": "true"},
"artifacts --json --repository web/app --collected"},
{"artifacts", map[string]any{"collected": "false"}, "artifacts --json"},
{"collect", map[string]any{}, "collect --json"},
{"collect", map[string]any{"most": "50"}, "collect --json --most 50"},
{"collect", map[string]any{"confirm": "true", "why": "the disk is full"},
"collect --json --confirm --why the disk is full"},
{"collect", map[string]any{"confirm": true, "why": "x", "most": "10"}, "collect --json --most 10 --confirm --why x"},
{"images", map[string]any{"node": "anchor"}, "images anchor --json"},
} {
argv, err := argvFor(c.verb, c.args)
if err != nil || strings.Join(argv, " ") != c.want {
t.Errorf("%s %v: %q %v, want %q", c.verb, c.args, argv, err, c.want)
}
}
for _, c := range []struct {
verb string
args map[string]any
}{
{"collect", map[string]any{"confirm": "true"}}, // a deletion without a why
{"collect", map[string]any{"why": "tidy"}}, // a why for a dry run, recorded nowhere
{"collect", map[string]any{"confirm": "yes", "why": "x"}}, // a switch is true or false
{"collect", map[string]any{"force": "true"}},
{"artifacts", map[string]any{"node": "anchor"}},
{"images", map[string]any{}},
{"images", map[string]any{"node": "anchor", "all": "true"}},
} {
if argv, err := argvFor(c.verb, c.args); err == nil {
t.Errorf("%s %v was composed as %q", c.verb, c.args, argv)
}
}
if repairingCommand([]string{"collect", "--json", "--confirm", "--why", "x"}) != "collect" {
t.Error("a real collect through the generic verb would go unrecorded")
}
if repairingCommand([]string{"collect", "--json"}) != "" {
t.Error("a dry run was taken for an act by hand")
}
if !personsDecision(link.HandAct{Verb: "collect"}) {
t.Error("a collect a person asked for would count toward a healer the mesh lacks")
}
}
func ref(module, artifact string, n int) string {
return fmt.Sprintf("%s%s/%s@sha256:%064x", catalogue.ArtifactStoreScheme, module, artifact, n)
}
func archiveRef(module, artifact string, n int) string {
return fmt.Sprintf("%s%s/%s/blobs/sha256:%064x", catalogue.ArtifactStoreScheme, module, artifact, n)
}
func TestArtifactsCountsEverythingAndListsTheCollectedOnlyWhenAsked(t *testing.T) {
states := []inventory.ArtifactState{
{Reference: ref("web", "app", 1), State: inventory.ArtifactCollected},
{Reference: ref("web", "app", 2), State: inventory.ArtifactEligible},
{Reference: archiveRef("web", "tools", 3), State: inventory.ArtifactKept,
Why: []string{inventory.KeptByRecentBuild, inventory.KeptByDefinition}},
{Reference: ref("db", "server", 4), State: inventory.ArtifactKept, Why: []string{inventory.KeptByDefinition}},
// Not the mesh's own store: never listed, never counted.
{Reference: "registry.invalid/x@sha256:" + strings.Repeat("a", 64), State: inventory.ArtifactKept,
Why: []string{inventory.KeptUnaddressable}},
}
a := artifactsOf(states, "", false)
if a.Counts != (artifactCounts{Kept: 2, Eligible: 1, Collected: 1}) {
t.Errorf("counts %+v", a.Counts)
}
if len(a.References) != 3 {
t.Fatalf("listed %d, want the kept and the eligible: %+v", len(a.References), a.References)
}
archive := a.References[1]
if archive.Kind != "archive" || archive.Repository != "web/tools" || archive.Digest != fmt.Sprintf("sha256:%064x", 3) ||
!slices.Equal(archive.Why, []string{"recent-build", "definition"}) {
t.Errorf("an archive read as %+v", archive)
}
if image := a.References[0]; image.Kind != "image" || image.Repository != "web/app" || image.State != "eligible" {
t.Errorf("an image read as %+v", image)
}
narrowed := artifactsOf(states, "web/app", true)
if narrowed.Counts != a.Counts || len(narrowed.References) != 2 {
t.Errorf("narrowed to one repository: %+v", narrowed)
}
raw, _ := json.Marshal(a.References[0])
if strings.Contains(string(raw), `"why"`) {
t.Errorf("an empty why is carried: %s", raw)
}
}
// fakeStore answers as a registry does for the sweep's questions, and records every request.
type fakeStore struct {
mu sync.Mutex
requests []string
refuse string // a DELETE whose path contains this is refused
}
func (f *fakeStore) serve(t *testing.T) artifacts.Store {
t.Helper()
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
f.mu.Lock()
f.requests = append(f.requests, r.Method+" "+r.URL.Path)
refuse := f.refuse
f.mu.Unlock()
switch {
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/blobs/"):
w.Header().Set("Content-Length", "10")
w.WriteHeader(http.StatusOK)
case r.Method == http.MethodHead:
w.WriteHeader(http.StatusNotFound)
case r.Method == http.MethodPost:
w.Header().Set("Location", "/upload/x?state=1")
w.WriteHeader(http.StatusAccepted)
case r.Method == http.MethodPut:
w.WriteHeader(http.StatusCreated)
case r.Method == http.MethodDelete && refuse != "" && strings.Contains(r.URL.Path, refuse):
w.WriteHeader(http.StatusInternalServerError)
case r.Method == http.MethodDelete:
w.WriteHeader(http.StatusAccepted)
default:
w.WriteHeader(http.StatusBadRequest)
}
}))
t.Cleanup(server.Close)
return artifacts.Store{Address: strings.TrimPrefix(server.URL, "http://")}
}
func (f *fakeStore) seen() []string {
f.mu.Lock()
defer f.mu.Unlock()
return slices.Clone(f.requests)
}
func TestADryRunCollectAsksOnlyAndChangesNothing(t *testing.T) {
f := &fakeStore{}
store := f.serve(t)
kept := []string{archiveRef("web", "tools", 1)}
eligible := []string{ref("web", "app", 2), archiveRef("web", "tools", 3)}
a := runCollect(context.Background(), nil, store, eligible, kept, false,
sweepBounds{most: 10, budget: 5 * time.Second})
if !a.DryRun || a.Eligible != 2 || !slices.Equal(a.WouldLetGo, eligible) || len(a.LetGo) != 0 {
t.Errorf("a dry run answered %+v", a)
}
if a.KeptArchives != 1 || a.Held != 0 || a.Unheld != 1 {
t.Errorf("the kept archive is unheld in the fake store, and was said as %+v", a)
}
for _, r := range f.seen() {
if !strings.HasPrefix(r, "HEAD ") {
t.Errorf("a dry run asked the store %s", r)
}
}
}
func TestCollectSaysWhatItDoesNotListAndWithNoStoreAsksNothing(t *testing.T) {
var many []string
for i := range mostListed + 5 {
many = append(many, ref("web", "app", i))
}
a := runCollect(context.Background(), nil, artifacts.Store{}, many, nil, false, sweepBounds{most: 1, budget: time.Second})
if len(a.WouldLetGo) != mostListed || a.NotListed != 5 || a.Left != len(many) || a.Stopped == "" {
t.Errorf("with no store and %d eligible: %d listed, %d not, left %d, stopped %q",
len(many), len(a.WouldLetGo), a.NotListed, a.Left, a.Stopped)
}
}
func TestARealCollectHoldsEveryKeptArchiveBeforeItLetsAnythingGo(t *testing.T) {
inv := inventory.ForTest(t)
f := &fakeStore{}
store := f.serve(t)
kept := []string{archiveRef("web", "tools", 1)}
eligible := []string{ref("web", "app", 2), ref("web", "app", 3), ref("web", "app", 4)}
madeBy(t, inv, eligible)
a := runCollect(t.Context(), inv, store, eligible, kept, true, sweepBounds{most: 2, budget: 5 * time.Second})
if a.DryRun || a.HoldersWritten != 1 || a.Held != 1 {
t.Errorf("the kept archive was not held first: %+v", a)
}
if !slices.Equal(a.LetGo, eligible[:2]) || a.Left != 1 || a.Stopped == "" {
t.Errorf("bounded at two: let go %v, left %d, stopped %q", a.LetGo, a.Left, a.Stopped)
}
requests := f.seen()
firstPut := slices.IndexFunc(requests, func(r string) bool { return strings.HasPrefix(r, "PUT ") && strings.Contains(r, "/manifests/") })
firstDelete := slices.IndexFunc(requests, func(r string) bool { return strings.HasPrefix(r, "DELETE ") })
if firstPut < 0 || firstDelete < 0 || firstPut > firstDelete {
t.Errorf("the holder was not put before the first delete: %v", requests)
}
// What was let go is recorded collected: the sweep no longer offers it.
left, err := inv.ToCollect(t.Context())
if err != nil {
t.Fatal(err)
}
if !slices.Equal(left, eligible[2:]) {
t.Errorf("after letting go of two, the records offer %v", left)
}
}
// madeBy records one successful build, of a module the mesh does not hold, that made these images:
// eligible, every one.
func madeBy(t *testing.T, inv *inventory.Inventory, references []string) {
t.Helper()
b := inventory.Build{ID: "b1", Repository: "https://forge.invalid/web.git", Module: "web", On: "a-build-machine",
Commit: "c0ffee"}
for i, r := range references {
b.Made = append(b.Made, inventory.Artifact{Name: fmt.Sprintf("app%d", i), Kind: "image", Reference: r})
}
if err := inv.RecordBuild(t.Context(), b); err != nil {
t.Fatal(err)
}
}
func TestARealCollectStopsAtTheStoresFirstRefusal(t *testing.T) {
inv := inventory.ForTest(t)
f := &fakeStore{refuse: fmt.Sprintf("%064x", 3)}
store := f.serve(t)
eligible := []string{ref("web", "app", 2), ref("web", "app", 3), ref("web", "app", 4)}
a := runCollect(t.Context(), inv, store, eligible, nil, true, sweepBounds{most: 10, budget: 5 * time.Second})
if !slices.Equal(a.LetGo, eligible[:1]) || a.Left != 2 || !strings.Contains(a.Stopped, "kept") {
t.Errorf("after a refusal: let go %v, left %d, stopped %q", a.LetGo, a.Left, a.Stopped)
}
}
func TestImagesAreEveryContainerImageOfTheDeclarationAndOfWhatWasSent(t *testing.T) {
body := []byte(`{"declaration":1,"resources":[
{"id":"web.server","type":"container","image":"store.invalid/web/server@sha256:aa"},
{"id":"web.collect","type":"container","image":"store.invalid/web/server@sha256:aa","schedule":"30 3 * * *"},
{"id":"db.server","type":"container","image":"postgres@sha256:bb"},
{"id":"web.bundle-tools","type":"archive","source":"http://store.invalid/v2/web/tools/blobs/sha256:cc"}]}`)
sent := []sentResource{
{ID: "web.server", Type: "container", Image: "store.invalid/web/server@sha256:99"},
{ID: "db.server", Type: "container", Image: "postgres@sha256:bb"},
{ID: "web.state", Type: "directory"},
}
a, err := imagesOf("anchor", body, sent, true)
if err != nil {
t.Fatal(err)
}
if !a.SentKnown || len(a.Images) != 3 {
t.Fatalf("answered %+v", a)
}
byImage := map[string]declaredImage{}
for _, d := range a.Images {
byImage[d.Image] = d
}
if d := byImage["store.invalid/web/server@sha256:aa"]; !slices.Equal(d.Resources, []string{"web.collect", "web.server"}) ||
!slices.Equal(d.In, []string{"declaration"}) {
t.Errorf("a scheduled step's image and its server's are one: %+v", d)
}
if d := byImage["store.invalid/web/server@sha256:99"]; !slices.Equal(d.In, []string{"sent"}) {
t.Errorf("the image last sent is kept beside the one to send: %+v", d)
}
if d := byImage["postgres@sha256:bb"]; !slices.Equal(d.In, []string{"declaration", "sent"}) {
t.Errorf("an image in both: %+v", d)
}
// A summary kept before it carried images cannot say what was sent.
old, err := imagesOf("anchor", body, []sentResource{{ID: "web.server", Type: "container"}}, true)
if err != nil || old.SentKnown {
t.Errorf("an old summary was read as knowing what was sent: %+v %v", old, err)
}
none, err := imagesOf("anchor", body, nil, false)
if err != nil || none.SentKnown || len(none.Images) != 2 {
t.Errorf("with nothing sent yet: %+v %v", none, err)
}
}
func TestTheSentSummaryKeepsAContainersImage(t *testing.T) {
summary, err := summarize([]byte(`{"resources":[{"id":"web.server","type":"container","image":"x@sha256:aa"},
{"id":"web.state","type":"directory","path":"/srv"}]}`))
if err != nil {
t.Fatal(err)
}
if summary[0].Image != "x@sha256:aa" || summary[1].Image != "" {
t.Errorf("summarized as %+v", summary)
}
}
+36
View File
@@ -573,6 +573,38 @@ func (a *verbArguments) commandLine() ([]string, error) {
return argv, nil
}
return []string{"cleanup", "list", "--json"}, nil
// novox/hq ADR 0251.
case "artifacts":
argv := []string{"artifacts", "--json"}
if r := str("repository"); r != "" {
argv = append(argv, "--repository", r)
}
if on("collected") {
argv = append(argv, "--collected")
}
return argv, nil
case "collect":
argv := []string{"collect", "--json"}
if m := str("most"); m != "" {
argv = append(argv, "--most", m)
}
why := str("why")
if !on("confirm") {
if why != "" {
return nil, errors.New("collect 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: letting go of artifacts is a hand act, which says why. Nothing was done", err)
}
return append(argv, "--confirm", "--why", why), nil
case "images":
if err := need("node"); err != nil {
return nil, err
}
return []string{"images", str("node"), "--json"}, nil
case "data":
argv := []string{"data", "--json"}
if m := str("machine"); m != "" {
@@ -768,6 +800,8 @@ func (a *verbArguments) commandLine() ([]string, error) {
// jsonVerbs are the verbs whose command speaks JSON, so the answer carries it as data as well.
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,
// The delivery's owner's verbs answer JSON where they read (plan, order, check, walks) — novox/hq ADR 0239.
"delivery": true}
@@ -791,6 +825,8 @@ func repairingCommand(argv []string) string {
return "retire " + argv[1]
case argv[0] == "cleanup" && len(argv) > 1 && argv[1] == "delete":
return "cleanup delete"
case argv[0] == "collect" && slices.Contains(argv, "--confirm"):
return "collect"
}
return ""
}
@@ -281,6 +281,9 @@ var accountedFlags = map[string]map[string]string{
"retire": {"json": "set by the verb: the answer is data"},
"cleanup": {"json": "set by the verb: the answer is data"},
"data": {"json": "set by the verb: the answer is data"},
"artifacts": {"json": "set by the verb: the answer is data"},
"collect": {"json": "set by the verb: the answer is data"},
"images": {"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"},
+7
View File
@@ -87,6 +87,10 @@ type sentResource struct {
// Volumes is a container's mounts as sent — paths, which are not secrets, and what a moved
// data directory is read from.
Volumes []string `json:"volumes,omitempty"`
// Image is a container's image as sent — a reference, not a secret — so the machine can be told
// which images its last declaration named, and keep them (novox/hq ADR 0251, `images`). Absent in
// a summary kept before it was.
Image string `json:"image,omitempty"`
}
// summarize is what is kept of a declaration body: its resources, in the order sent.
@@ -105,6 +109,9 @@ func summarize(body []byte) ([]sentResource, error) {
s.Fields[k] = digestValue(v)
}
if s.Type == "container" {
if image, ok := r["image"].(string); ok {
s.Image = image
}
if vols, ok := r["volumes"].([]any); ok {
for _, v := range vols {
s.Volumes = append(s.Volumes, fmt.Sprint(v))
+28
View File
@@ -435,6 +435,34 @@ var ControllerVerbs = []Verb{
"behind": "\"true\", instead of a repository: build every module the mesh holds whose source has moved " +
"past the commit it was built from, bases asked first",
}, nil, "behind")},
// What the records say the registries may keep, asked on demand (novox/hq ADR 0251, to-be 51).
{Name: "artifacts", Description: "Every artifact the mesh recorded making that lives in its artifact store, " +
"each with its repository, digest, kind (image or archive) and state: kept — and why: a definition names " +
"it, or one of the five most recent builds of a module the mesh holds — or eligible: made, kept for no " +
"reason, not yet let go. Collected ones are counted, and listed only with collected. Reads only " +
"(novox/hq ADR 0251).",
Input: schema(map[string]string{
"repository": "only this repository's references, as <module>/<artifact> (the counts stay whole)",
"collected": "\"true\": list the collected references too (they only grow; counted always)",
}, nil, "collected")},
{Name: "collect", Description: "The artifact store's sweep, now: what the mesh made and keeps for no reason " +
"is let go of. Without confirm a dry run — every kept archive asked about with HEADs, what would be let go " +
"listed, nothing held or deleted. With confirm and why: every kept archive held first, then each eligible " +
"artifact let go of, oldest first, and recorded collected — a hand act, recorded with its why. Bounded by " +
"most (500 by default, at most 5000) and 45 seconds; what is left is said. Bytes are reclaimed by the " +
"store's nightly collector. Never touches a digest the mesh did not record (novox/hq ADR 0189, ADR 0251).",
Input: schema(map[string]string{
"confirm": "\"true\": let go of what is eligible (needs why); without it nothing is done",
"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: "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 " +
"its images. Reads only (novox/hq ADR 0251).",
Input: schema(map[string]string{
"node": "the machine's name",
}, []string{"node"})},
}
// schema is a JSON schema for an object of string properties, which is every argument the verbs
+66
View File
@@ -0,0 +1,66 @@
package inventory
import (
"context"
"fmt"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// **Every artifact made has one state, and a kept one says why** (novox/hq ADR 0251): what the
// store's own tools set beside what the store holds, so an operator reads why something stays.
func TestEveryArtifactMadeHasAStateAndAKeptOneSaysWhy(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
var made []string
for i := 1; i <= 8; i++ {
made = append(made, built(t, inv, fmt.Sprintf("b%02d", i), "web", i))
}
// The definition names the oldest and the newest: the newest is kept for both reasons.
m := catalogue.Manifest{Module: "web", Version: "1", Resources: []map[string]any{
{"id": "app", "type": "container", "name": "web", "image": made[0]},
{"id": "next", "type": "container", "name": "web-next", "image": made[7]},
}}
if err := inv.RegisterModule(ctx, m, Source{Repository: "https://forge.invalid/web.git"}); err != nil {
t.Fatal(err)
}
if err := inv.MarkCollected(ctx, []string{made[1]}); err != nil {
t.Fatal(err)
}
states, err := inv.Artifacts(ctx)
if err != nil {
t.Fatal(err)
}
if len(states) != len(made) {
t.Fatalf("%d states for %d artifacts made", len(states), len(made))
}
want := []struct {
state string
why []string
}{
{ArtifactKept, []string{KeptByDefinition}},
{ArtifactCollected, nil},
{ArtifactEligible, nil},
{ArtifactKept, []string{KeptByRecentBuild}},
{ArtifactKept, []string{KeptByRecentBuild}},
{ArtifactKept, []string{KeptByRecentBuild}},
{ArtifactKept, []string{KeptByRecentBuild}},
{ArtifactKept, []string{KeptByRecentBuild, KeptByDefinition}},
}
for i, s := range states {
if s.Reference != made[i] {
t.Fatalf("state %d is of %s, want %s: the oldest first", i, s.Reference, made[i])
}
if s.State != want[i].state || fmt.Sprint(s.Why) != fmt.Sprint(want[i].why) {
t.Errorf("%s: %s %v, want %s %v", s.Reference, s.State, s.Why, want[i].state, want[i].why)
}
}
collect, err := inv.ToCollect(ctx)
if err != nil {
t.Fatal(err)
}
if len(collect) != 1 || collect[0] != made[2] {
t.Errorf("ToCollect is %v, want only the eligible %s", collect, made[2])
}
}
+88 -50
View File
@@ -27,7 +27,39 @@ import (
// release that turns out wrong can be taken.
const KeptBuilds = 5
// ToCollect is every artifact the mesh made, no longer keeps, and has not already collected.
// Why an artifact is kept: the words `artifacts` answers (novox/hq ADR 0251).
const (
// KeptByDefinition: a definition the mesh holds names it, so it could be handed to a machine now.
KeptByDefinition = "definition"
// KeptByRecentBuild: it is an artifact of one of the KeptBuilds most recent successful builds of a
// module the mesh still holds — somewhere a release that turns out wrong can go back to.
KeptByRecentBuild = "recent-build"
// KeptUnaddressable: it names no digest, so nothing here can speak for it, and it is kept rather
// than guessed about.
KeptUnaddressable = "unaddressable"
)
// The states an artifact the mesh recorded making is in (novox/hq ADR 0251).
const (
ArtifactKept = "kept"
ArtifactEligible = "eligible"
ArtifactCollected = "collected"
)
// ArtifactState is one artifact the mesh recorded making, and what the records say of it now.
type ArtifactState struct {
Reference string
State string
// Why is the reasons a kept artifact is kept; empty for any other.
Why []string
}
// Artifacts is every artifact any successful build recorded, oldest first, each with its state:
// kept (and why), eligible to let go, or already collected (novox/hq ADR 0189, ADR 0251).
//
// **One reading of the records for every question asked of them**: the sweep's ToCollect, the
// operator's `artifacts` and the store's own tools read the same states, so what one says may go is
// what the others say is eligible.
//
// Three reasons an artifact stays, and nothing else is a reason:
//
@@ -38,57 +70,62 @@ const KeptBuilds = 5
// nothing beyond what a held definition names (novox/hq issue 253);
// - it was already collected, in which case there is nothing left to do.
//
// Returned in a stated order so two runs over the same records ask for the same things in the
// same sequence, which is what makes a failed sweep safe to simply run again.
func (i *Inventory) ToCollect(ctx context.Context) ([]string, error) {
keep, err := i.keptReferences(ctx)
// In a stated order, so two runs over the same records ask for the same things in the same
// sequence, which is what makes a failed sweep safe to simply run again.
func (i *Inventory) Artifacts(ctx context.Context) ([]ArtifactState, error) {
keep, err := i.KeepSet(ctx)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
// Every artifact of every successful build, oldest first, minus what has already been
// collected. A failed build published nothing, so it names nothing to remove.
`select b.made
from build b
where b.failed = '' and b.module is not null and b.module <> ''
order by b.at asc, b.id asc`)
all, err := i.everyReferenceMade(ctx)
if err != nil {
return nil, err
}
defer rows.Close()
collected, err := i.alreadyCollected(ctx)
if err != nil {
return nil, err
}
seen := map[string]bool{}
var out []string
for rows.Next() {
var raw []byte
if err := rows.Scan(&raw); err != nil {
return nil, err
}
var made []Artifact
if err := json.Unmarshal(raw, &made); err != nil {
// One unreadable record must not stop the rest being collected — and an artifact this
// row named is simply not offered, which errs toward keeping.
continue
}
for _, a := range made {
reference := asRecorded(a.Reference)
if reference == "" || keep[reference] || collected[reference] || seen[reference] {
continue
}
seen[reference] = true
out = append(out, reference)
out := make([]ArtifactState, 0, len(all))
for _, reference := range all {
switch {
case len(keep[reference]) > 0:
out = append(out, ArtifactState{Reference: reference, State: ArtifactKept, Why: keep[reference]})
case collected[reference]:
out = append(out, ArtifactState{Reference: reference, State: ArtifactCollected})
default:
out = append(out, ArtifactState{Reference: reference, State: ArtifactEligible})
}
}
return out, rows.Err()
return out, nil
}
// keptReferences is every artifact reference the mesh still keeps, for either of the two reasons.
func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error) {
keep := map[string]bool{}
// ToCollect is every artifact the mesh made, no longer keeps, and has not already collected —
// oldest first (Artifacts says why each other one stays).
func (i *Inventory) ToCollect(ctx context.Context) ([]string, error) {
states, err := i.Artifacts(ctx)
if err != nil {
return nil, err
}
var out []string
for _, s := range states {
if s.State == ArtifactEligible {
out = append(out, s.Reference)
}
}
return out, nil
}
// KeepSet is every artifact reference the mesh still keeps, each with the reasons it is kept.
func (i *Inventory) KeepSet(ctx context.Context) (map[string][]string, error) {
keep := map[string][]string{}
because := func(reference, why string) {
for _, w := range keep[reference] {
if w == why {
return
}
}
keep[reference] = append(keep[reference], why)
}
// **Whatever a definition the mesh holds names.** Read as text rather than by walking the
// resource shapes: a reference may be a container's image, a bundle's source, or a field some
@@ -140,7 +177,7 @@ func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error)
}
for _, a := range made {
if reference := asRecorded(a.Reference); reference != "" {
keep[reference] = true
because(reference, KeptByRecentBuild)
}
}
}
@@ -148,28 +185,28 @@ func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error)
return nil, err
}
// And anything a manifest mentions. Done after the recent set so the scan runs over the
// candidates rather than over every reference ever recorded: a manifest holds a reference
// composed with the store's address or kept bare, so the search is for the digest within it.
// 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
// search is for the digest within it.
if len(named) > 0 {
all, err := i.everyReferenceMade(ctx)
if err != nil {
return nil, err
}
for _, reference := range all {
if keep[reference] {
continue
}
digest := digestIn(reference)
if digest == "" {
// Not something the store holds by digest; nothing here can speak for it, so it
// is kept rather than guessed about.
keep[reference] = true
if len(keep[reference]) == 0 {
because(reference, KeptUnaddressable)
}
continue
}
for _, text := range named {
if strings.Contains(text, digest) {
keep[reference] = true
because(reference, KeptByDefinition)
break
}
}
@@ -186,7 +223,7 @@ func (i *Inventory) keptReferences(ctx context.Context) (map[string]bool, error)
// reference that is kept because nothing here can speak for it is not one the store can be asked
// about.
func (i *Inventory) KeptArchives(ctx context.Context) ([]string, error) {
keep, err := i.keptReferences(ctx)
keep, err := i.KeepSet(ctx)
if err != nil {
return nil, err
}
@@ -200,10 +237,11 @@ func (i *Inventory) KeptArchives(ctx context.Context) ([]string, error) {
return out, nil
}
// everyReferenceMade is every artifact reference any successful build recorded.
// 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,
`select made from build where failed = '' and module is not null and module <> ''`)
`select made from build where failed = '' and module is not null and module <> ''
order by at asc, id asc`)
if err != nil {
return nil, err
}
+4 -1
View File
@@ -72,7 +72,10 @@
"retire",
"cleanup",
"data",
"build"
"build",
"artifacts",
"collect",
"images"
],
"resources": [
{