Files
mesh-controller/internal/inventory/collection.go
T
jochen 9907df6530 Record the bases a build copies, keep them by the builds that stood on them, and copy each image once
A copied base was named only in what a build stood on, and nowhere when the build failed,
so the store's sweep could never let one go (hq issue 321). One repository per upstream
image stops each module asking the public registry for the same image again, and letting
an index go now takes its own platform manifests, which otherwise kept every byte. A
person can record the copies no record names through the new mirrors verb (hq ADR 0257).

The forge test fix is the same commit as on feat/plain-notifications: main fails without it.
2026-10-08 15:47:42 +02:00

442 lines
16 KiB
Go

package inventory
import (
"context"
"encoding/json"
"sort"
"strings"
"github.com/novox/mesh-controller/internal/catalogue"
)
// What the artifact store keeps, and what it may let go (novox/hq ADR 0189, issue 108).
//
// The store has never collected anything: every build pushes another layer set and nothing has
// ever removed one. The registry's own answer — collect what no tag names — is wrong here, because
// the mesh pushes each artifact under one moving tag and pins machines by digest, so every build
// but the newest is untagged and some machine may still be running it.
//
// **So the mesh decides, from its own records, and it never has to look in the store to do it.**
// It has never put anything there it did not record, which means every digest it could remove is
// already in a build row. A digest the mesh did not record making is therefore never named here —
// not as a safety margin but as the rule restated, and it is what keeps the sweep away from the
// images genesis pushed before any record existed (04-ISSUES/102, F4).
// KeptBuilds is how many successful builds of each module keep their artifacts, counting the
// newest. The newest is what the mesh hands a machine now; the four behind it are how far back a
// release that turns out wrong can be taken.
const KeptBuilds = 5
// 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"
// KeptStoodOn: it is a base the mesh copied in, and one of the KeptBuilds most recent successful
// builds of a module the mesh still holds stood on it (novox/hq ADR 0257). A copy has no builds of
// its own; it is kept by the builds that used it.
KeptStoodOn = "stood-on"
)
// The states an artifact the mesh recorded making is in (novox/hq ADR 0251).
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:
//
// - **a definition names it** — the reference appears in a module's recorded manifest, which is
// what the mesh would hand a machine now. No age limit: this is the floor;
// - **the mesh can still go back to it** — it is an artifact of one of the KeptBuilds most
// recent successful builds of a module the mesh still holds. A forgotten module keeps
// nothing beyond what a held definition names (novox/hq issue 253);
// - it was already collected, in which case there is nothing left to do.
//
// 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
}
all, err := i.everyReferenceRecorded(ctx)
if err != nil {
return nil, err
}
collected, err := i.alreadyCollected(ctx)
if err != nil {
return nil, err
}
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, nil
}
// 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
// later kind of resource grows, and what matters is only whether the mesh could hand this
// string to a machine. A manifest that mentions it is a manifest that might.
manifests, err := i.store.Pool().Query(ctx, `select manifest::text from module where manifest is not null`)
if err != nil {
return nil, err
}
defer manifests.Close()
var named []string
for manifests.Next() {
var text string
if err := manifests.Scan(&text); err != nil {
return nil, err
}
named = append(named, text)
}
if err := manifests.Err(); err != nil {
return nil, err
}
// The KeptBuilds most recent successful builds of each module the mesh still holds, whole.
//
// **Only a module the mesh still holds can be gone back to** (novox/hq issue 253). "Somewhere
// to return to" is a reason about a module's releases; a module that has been forgotten has
// no releases left to return between, and its build rows stay only as history. Without the
// join every module ever built kept five builds' artifacts for ever — and once the store's
// collector runs for real, what the keep set says is what the disk holds.
recent, err := i.store.Pool().Query(ctx,
`select made, built_against from (
select b.made, b.built_against,
row_number() over (partition by b.module order by b.at desc, b.id desc) as back
from build b
join module m on m.name = b.module
where b.failed = '' and b.module is not null and b.module <> ''
) ranked where back <= $1`, KeptBuilds)
if err != nil {
return nil, err
}
defer recent.Close()
var stoodOn []string
for recent.Next() {
var raw, against []byte
if err := recent.Scan(&raw, &against); err != nil {
return nil, err
}
var made []Artifact
if err := json.Unmarshal(raw, &made); err == nil {
for _, a := range made {
if reference := asRecorded(a.Reference); reference != "" {
because(reference, KeptByRecentBuild)
}
}
}
var bases []string
if err := json.Unmarshal(against, &bases); err == nil {
for _, b := range bases {
if reference := asRecorded(b); reference != "" {
stoodOn = append(stoodOn, reference)
}
}
}
}
if err := recent.Err(); err != nil {
return nil, err
}
// **A copied base is kept by the kept builds that stood on it, and by nothing else** (novox/hq
// ADR 0257). It has no builds of its own to count five of, and no machine is ever handed a base,
// so a definition is never its reason: a definition names the image upstream, by the digest every
// copy of it shares, and read as a reason it would keep every module's copy for as long as any
// module stood on that image. A copy already let go of stays let go of.
mirrored, err := i.mirroredSet(ctx)
if err != nil {
return nil, err
}
if len(mirrored) > 0 {
collected, err := i.alreadyCollected(ctx)
if err != nil {
return nil, err
}
for _, reference := range stoodOn {
if mirrored[reference] && !collected[reference] {
because(reference, KeptStoodOn)
}
}
}
// And anything a manifest mentions — every reference ever made, the recent ones included, so a
// kept artifact carries each of its reasons and the operator reads why, not only that (novox/hq
// ADR 0251). A manifest holds a reference composed with the store's address or kept bare, so the
// 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 mirrored[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.
if len(keep[reference]) == 0 {
because(reference, KeptUnaddressable)
}
continue
}
for _, text := range named {
if strings.Contains(text, digest) {
because(reference, KeptByDefinition)
break
}
}
}
}
return keep, nil
}
// KeptArchives is every archive the mesh keeps, in a stated order: the references the sweep must
// hold by a manifest before it lets anything go, and the ones an operator needs to read as all
// held before the store's collector is let loose (novox/hq issue 253, ADR 0189).
//
// Only references into the mesh's own store, and only blobs: an image is its own manifest, and a
// 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.KeepSet(ctx)
if err != nil {
return nil, err
}
var out []string
for reference := range keep {
if path, ours := catalogue.InArtifactStore(reference); ours && strings.Contains(path, "/blobs/sha256:") {
out = append(out, reference)
}
}
sort.Strings(out)
return out, nil
}
// KeptImages is every image the mesh keeps, in a stated order: the manifests whose platform manifests
// a sweep spares when it lets an index of the same repository go (novox/hq ADR 0257).
func (i *Inventory) KeptImages(ctx context.Context) ([]string, error) {
keep, err := i.KeepSet(ctx)
if err != nil {
return nil, err
}
var out []string
for reference := range keep {
if path, ours := catalogue.InArtifactStore(reference); ours && strings.Contains(path, "@sha256:") {
out = append(out, reference)
}
}
sort.Strings(out)
return out, nil
}
// everyReferenceRecorded is every artifact the mesh recorded putting in its store: what each
// successful build made, and every base a build — or a person — recorded copying in (novox/hq ADR
// 0257). Oldest first, a copy placed at the moment it was recorded.
func (i *Inventory) everyReferenceRecorded(ctx context.Context) ([]string, error) {
made, err := i.everyReferenceMade(ctx)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx, `select reference from artifact_mirrored order by at asc, reference asc`)
if err != nil {
return nil, err
}
defer rows.Close()
seen := map[string]bool{}
for _, r := range made {
seen[r] = true
}
out := made
for rows.Next() {
var reference string
if err := rows.Scan(&reference); err != nil {
return nil, err
}
if reference = asRecorded(reference); reference != "" && !seen[reference] {
seen[reference] = true
out = append(out, reference)
}
}
return out, rows.Err()
}
// mirroredSet is every base the records say the mesh copied in, as recorded.
func (i *Inventory) mirroredSet(ctx context.Context) (map[string]bool, error) {
rows, err := i.store.Pool().Query(ctx, `select reference from artifact_mirrored`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]bool{}
for rows.Next() {
var reference string
if err := rows.Scan(&reference); err != nil {
return nil, err
}
out[asRecorded(reference)] = true
}
return out, rows.Err()
}
// everyReferenceMade is every artifact reference any successful build recorded, oldest build first.
func (i *Inventory) everyReferenceMade(ctx context.Context) ([]string, error) {
rows, err := i.store.Pool().Query(ctx,
`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
}
defer rows.Close()
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 {
continue
}
for _, a := range made {
reference := asRecorded(a.Reference)
if reference == "" || seen[reference] {
continue
}
seen[reference] = true
out = append(out, reference)
}
}
return out, rows.Err()
}
// asRecorded is an artifact reference in the one vocabulary the sweep speaks (novox/hq issue 226).
//
// **Every reference here came from a build record, so every one of them is the mesh's own.** That
// is what makes it safe to normalise: references kept before the store's address stopped being
// written are `<host>:<port>/<path>@sha256:…` (04-ISSUES/102), and `Recorded` reads those as the
// `artifact-store://` references the rest of the mesh uses. Done here rather than when the store
// is asked, because `Recorded` cannot tell one registry host from another — only the provenance
// can, and the provenance is here.
//
// The oldest artifacts are exactly the ones recorded the old way, and exactly the ones a
// sweep reaches first. Untranslated, the first of them ended every sweep.
func asRecorded(reference string) string {
if reference == "" {
return ""
}
return catalogue.Recorded(reference)
}
// digestIn is the `sha256:<hex>` a reference names, empty when it names none.
func digestIn(reference string) string {
for _, marker := range []string{"@sha256:", "/sha256:"} {
if _, after, ok := strings.Cut(reference, marker); ok {
return "sha256:" + after
}
}
return ""
}
// alreadyCollected is what the store has already been asked to let go.
func (i *Inventory) alreadyCollected(ctx context.Context) (map[string]bool, error) {
rows, err := i.store.Pool().Query(ctx, `select reference from artifact_collected`)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]bool{}
for rows.Next() {
var reference string
if err := rows.Scan(&reference); err != nil {
return nil, err
}
out[reference] = true
}
return out, rows.Err()
}
// MarkCollected records that the store no longer holds these.
//
// **A store that answered "not found" is recorded too.** The outcome wanted is that the artifact
// is gone, and it is; retrying it every sweep for ever is the failure this table exists to
// prevent. Only a store that could not be reached, or refused, leaves a reference unmarked — and
// then the next sweep asks again, which is what should happen.
func (i *Inventory) MarkCollected(ctx context.Context, references []string) error {
for _, reference := range references {
if _, err := i.store.Pool().Exec(ctx,
`insert into artifact_collected (reference) values ($1) on conflict (reference) do nothing`,
reference); err != nil {
return err
}
}
return nil
}
// Mirrored is every base the records say the mesh copied into its store (novox/hq ADR 0257).
func (i *Inventory) Mirrored(ctx context.Context) (map[string]bool, error) {
return i.mirroredSet(ctx)
}