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).
444 lines
14 KiB
Go
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
|
|
}
|