diff --git a/cmd/mesh-controller/collect.go b/cmd/mesh-controller/collect.go index d30344e1..f9c1362c 100644 --- a/cmd/mesh-controller/collect.go +++ b/cmd/mesh-controller/collect.go @@ -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) } } diff --git a/cmd/mesh-controller/handacts.go b/cmd/mesh-controller/handacts.go index 9901927c..e127d58e 100644 --- a/cmd/mesh-controller/handacts.go +++ b/cmd/mesh-controller/handacts.go @@ -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 " + diff --git a/cmd/mesh-controller/main.go b/cmd/mesh-controller/main.go index a25a1127..71c291fb 100644 --- a/cmd/mesh-controller/main.go +++ b/cmd/mesh-controller/main.go @@ -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 --why the provider deletes that one retired consumer cleanup delete --older-than --why [--confirm] list those older; delete only with --confirm + artifacts [--json] [--repository ] [--collected] every artifact the mesh made: kept and why, eligible, collected (ADR 0251) + collect [--json] [--most ] what the store's sweep would let go of, now; nothing is done + collect --confirm --why [--most ] run the sweep now: hold every kept archive, let go of the eligible + images [--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 rename a seat; its former name still resolves (ADR 0122) seat --to / hand a seat to that assignment as one act; never empty in between (ADR 0131) diff --git a/cmd/mesh-controller/registry_verbs.go b/cmd/mesh-controller/registry_verbs.go new file mode 100644 index 00000000..ab4dd994 --- /dev/null +++ b/cmd/mesh-controller/registry_verbs.go @@ -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 (/); 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 /] [--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 ] [--most ]") + } + 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 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 [--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 +} diff --git a/cmd/mesh-controller/registry_verbs_test.go b/cmd/mesh-controller/registry_verbs_test.go new file mode 100644 index 00000000..f1d08cac --- /dev/null +++ b/cmd/mesh-controller/registry_verbs_test.go @@ -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) + } +} diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 1eb0ccfa..d299a805 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -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 "" } diff --git a/cmd/mesh-controller/seatverbs_schema_test.go b/cmd/mesh-controller/seatverbs_schema_test.go index b0f2243e..e9a3cf4b 100644 --- a/cmd/mesh-controller/seatverbs_schema_test.go +++ b/cmd/mesh-controller/seatverbs_schema_test.go @@ -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"}, diff --git a/cmd/mesh-controller/unseen.go b/cmd/mesh-controller/unseen.go index 543a96e6..1ed42bcd 100644 --- a/cmd/mesh-controller/unseen.go +++ b/cmd/mesh-controller/unseen.go @@ -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)) diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index a6601ff6..1e637ace 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -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 / (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 diff --git a/internal/inventory/artifacts_test.go b/internal/inventory/artifacts_test.go new file mode 100644 index 00000000..5deef9cc --- /dev/null +++ b/internal/inventory/artifacts_test.go @@ -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]) + } +} diff --git a/internal/inventory/collection.go b/internal/inventory/collection.go index 684287ab..5931da42 100644 --- a/internal/inventory/collection.go +++ b/internal/inventory/collection.go @@ -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 } diff --git a/module.json b/module.json index 0ee497d0..b61e333f 100644 --- a/module.json +++ b/module.json @@ -72,7 +72,10 @@ "retire", "cleanup", "data", - "build" + "build", + "artifacts", + "collect", + "images" ], "resources": [ {