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

444 lines
14 KiB
Go

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, platforms: true}
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
}