package inventory import ( "context" "encoding/json" "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 // ToCollect is every artifact the mesh made, no longer keeps, and has not already collected. // // 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 its module; // - 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) 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`) 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) } } return out, rows.Err() } // 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{} // **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, whole. recent, err := i.store.Pool().Query(ctx, `select made from ( select made, row_number() over (partition by module order by at desc, id desc) as back from build where failed = '' and module is not null and module <> '' ) ranked where back <= $1`, KeptBuilds) if err != nil { return nil, err } defer recent.Close() for recent.Next() { var raw []byte if err := recent.Scan(&raw); err != nil { return nil, err } var made []Artifact if err := json.Unmarshal(raw, &made); err != nil { continue } for _, a := range made { if reference := asRecorded(a.Reference); reference != "" { keep[reference] = true } } } if err := recent.Err(); err != nil { 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. 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 continue } for _, text := range named { if strings.Contains(text, digest) { keep[reference] = true break } } } } return keep, nil } // everyReferenceMade is every artifact reference any successful build recorded. 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 <> ''`) 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 `:/@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:` 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 }