Files
mesh-controller/cmd/mesh-controller/collect.go
T
jochen c6e372896b
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer delivery to the same trunk took over its walk
Leave an eligible index for a confirmed collect instead of letting it go alone
Let go of alone after a build, an index's platforms stayed for ever under a record that
said collected, so no later collect could reach them (re-review of #144).
2026-10-08 15:47:42 +02:00

277 lines
12 KiB
Go

package main
import (
"context"
"errors"
"fmt"
"os"
"time"
"github.com/novox/mesh-controller/internal/artifacts"
"github.com/novox/mesh-controller/internal/inventory"
)
// Letting the artifact store go of what the mesh no longer keeps (novox/hq ADR 0189, issue 108).
//
// **Run where the records change.** A build is the moment new bytes landed in the store and the
// 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
// platforms is whether an index is let go of with its own platform manifests (novox/hq ADR 0257).
// Only a sweep a person confirmed, after its dry run listed what stays and names a platform: the
// records cannot see an unrecorded index, or a copy in progress, naming the same platform, and the
// store's tools can.
platforms bool
}
// 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
// Indexes is how many eligible indexes a sweep a person did not confirm left for one that is: an
// index goes only with its platform manifests (novox/hq ADR 0257 §4).
Indexes 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.
// **An index goes with its platform manifests only when a person confirmed it** (novox/hq ADR
// 0257), and a kept index keeps its own: which platform manifests the kept indexes name is read
// from the store for every repository the sweep will touch, before the first delete. A single
// read that fails stops the sweep before it deletes anything, as a kept archive that cannot be
// held does. The sweep after a build leaves an eligible index for such a collect, below.
store.Spare = nil
if bounds.platforms {
if inv == nil {
r.Stopped = "no records to read which platform manifests a kept index names, so nothing was let go"
r.Left = len(references)
return r
}
images, err := inv.KeptImages(ctx)
if err == nil {
store.Spare, err = store.SpareKeptIndexes(within, images, references)
}
if err != nil {
r.Stopped = fmt.Sprintf("which platform manifests a kept index names could not be read, so nothing was let go: %v", err)
r.Left = len(references)
return r
}
}
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
}
if !bounds.platforms {
// **An index waits for a confirmed collect** (novox/hq ADR 0257 §4). Let go of alone,
// its platform manifests would stay in the store and the record would say collected, so
// no later sweep would offer it again and they would stay for ever. Left eligible, the
// next collect a person confirms takes it with them.
index, err := store.IsIndex(within, reference)
if err != nil && !errors.Is(err, artifacts.Gone) && !errors.Is(err, artifacts.ErrNotOurs) {
r.Stopped = fmt.Sprintf("the artifact store could not say whether %s is an index, so nothing more was asked of it: %v", reference, err)
r.Left = len(references) - i
break
}
if index {
r.Indexes++
continue
}
}
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 — 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 {
fmt.Fprintf(os.Stderr, "could not work out what the artifact store may let go of: %v\n", err)
return
}
kept, err := inv.KeptArchives(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "could not work out which archives the artifact store keeps: %v\n", err)
return
}
if len(references) == 0 && len(kept) == 0 {
return
}
shelf, err := inv.Catalogue(ctx)
if err != nil {
fmt.Fprintf(os.Stderr, "could not read the catalogue to find the artifact store: %v\n", err)
return
}
// As the mesh reaches it from the network. Empty means the store is not on the network — on a
// mesh being raised it is not yet, and there the store holds one build of anything and has
// nothing to collect.
address, err := artifactStoreAddress(ctx, inv, shelf, "")
if err != nil || address == "" {
if err != nil {
fmt.Fprintf(os.Stderr, "could not find the artifact store to collect from: %v\n", err)
}
return
}
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 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", r.Missing)
}
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))
}
if r.Stopped != "" && !r.Bounded {
fmt.Fprintln(os.Stderr, r.Stopped)
}
if r.Held && r.Left > 0 {
fmt.Fprintf(os.Stderr, "%d more to collect; the next build asks again\n", r.Left)
}
if r.Skipped > 0 {
fmt.Fprintf(os.Stderr, "%d artifact(s) the sweep will not address were skipped\n", r.Skipped)
}
if r.Indexes > 0 {
fmt.Fprintf(os.Stderr, "%d eligible index(es) wait for a collect a person confirms, which takes their platforms too\n", r.Indexes)
}
}
// holdKept holds every kept archive by its manifest, stopping at the first refusal by the store.
// Answers how many holders it wrote and how many kept archives the store does not have.
//
// A missing archive is counted rather than fatal: there is nothing to hold, and that is a fact
// for an operator to read (`collection`), not a reason to stop collecting what is not kept. A
// reference the store cannot be asked about is skipped as the deletion loop skips one
// (novox/hq issue 226).
func holdKept(ctx context.Context, store artifacts.Store, kept []string) (wrote, missing int, err error) {
for _, reference := range kept {
if err := ctx.Err(); err != nil {
return wrote, missing, fmt.Errorf("ran out of time before %s: %w", reference, err)
}
did, err := store.Hold(ctx, reference)
switch {
case err == nil:
if did {
wrote++
}
case errors.Is(err, artifacts.Gone):
missing++
case errors.Is(err, artifacts.ErrNotOurs):
default:
return wrote, missing, fmt.Errorf("holding %s: %w", reference, err)
}
}
return wrote, missing, nil
}
// mostPerSweep is how many artifacts one sweep will ask about. Enough that a mesh building
// several times a day converges within days of this landing; small enough that no single build
// waits on the whole backlog.
const mostPerSweep = 200
// sweepBudget is the longest a sweep will keep a build waiting.
const sweepBudget = 60 * time.Second