diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 5b969b2..d169034 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -533,6 +533,10 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu if err := inv.RegisterModule(ctx, manifest, recorded); err != nil { return manifest, kept, err } + // The keep set just moved, and new bytes just landed (novox/hq ADR 0189). Asked here rather + // than on a timer of its own: this is the only moment either is true. Never fatal — the build + // worked and the module is registered. + collect(ctx, inv) return manifest, kept, nil } diff --git a/cmd/mesh-controller/collect.go b/cmd/mesh-controller/collect.go new file mode 100644 index 0000000..c90e98a --- /dev/null +++ b/cmd/mesh-controller/collect.go @@ -0,0 +1,84 @@ +package main + +import ( + "context" + "errors" + "fmt" + "os" + + "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. +// +// 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. + +// 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. +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 + } + if len(references) == 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 + } + + store := artifacts.Store{Address: address} + var done []string + var refused int + for _, reference := range references { + switch err := store.LetGo(ctx, reference); { + case 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) + default: + refused++ + if refused == 1 { + // Once per sweep. A store that refuses one refuses all of them, and a hundred + // identical lines would bury the reason. + fmt.Fprintf(os.Stderr, "the artifact store kept %s: %v\n", reference, err) + } + } + } + if len(done) > 0 { + if err := inv.MarkCollected(ctx, done); err != nil { + // Said, and that is all: the artifacts are gone either way, and the only cost of an + // unrecorded collection is that the next sweep asks about them again. + 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 refused > 0 { + fmt.Fprintf(os.Stderr, "%d artifact(s) were not collected; the next build asks again\n", refused) + } +} diff --git a/internal/artifacts/store.go b/internal/artifacts/store.go new file mode 100644 index 0000000..02b79d8 --- /dev/null +++ b/internal/artifacts/store.go @@ -0,0 +1,99 @@ +// Package artifacts speaks to the mesh's artifact store over its own door. +// +// Only what the mesh needs that nothing else does: letting go of something it put there +// (novox/hq ADR 0189, issue 108). Pushing is the builder's, through the container runtime; reading +// is every machine's, through its runtime. This is the one operation that belongs to the thing +// holding the records, because it is the only one that is a decision rather than a transfer. +package artifacts + +import ( + "context" + "fmt" + "net/http" + "strings" + "time" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// Store is the artifact store at an address, as this machine reaches it. +type Store struct { + // Address is `host:port` — the store as the caller reaches it now, composed and never + // recorded (novox/hq 04-ISSUES/102). + Address string + // HTTP is the client used; nil is a client with a modest timeout. + HTTP *http.Client +} + +// Gone is the answer when the store does not hold it: the outcome wanted, already true. +var Gone = fmt.Errorf("the store does not hold it") + +// LetGo asks the store to drop one artifact the mesh recorded making. +// +// Takes a reference as the mesh records it — `artifact-store:///@sha256:…` for +// an image, `…/blobs/sha256:…` for an archive — because that is the identity every record uses, +// and composes the address here at the moment of use. +// +// Returns Gone when the store answers that it does not have it. That is not a failure: the sweep +// wants the artifact absent, and it is. It is distinguished from success only so a caller can say +// which of the two happened. +func (s Store) LetGo(ctx context.Context, reference string) error { + path, kept := catalogue.InArtifactStore(reference) + if !kept { + // Nothing the mesh put in its own store. Refused rather than attempted: composing a + // delete for a reference of unknown shape is how a sweep reaches something that is not + // the mesh's. + return fmt.Errorf("%s is not a reference into the mesh's artifact store", reference) + } + if s.Address == "" { + return fmt.Errorf("this mesh has no artifact store on its network to ask about %s", reference) + } + repository, kind, digest, err := split(path) + if err != nil { + return err + } + url := "http://" + s.Address + "/v2/" + repository + "/" + kind + "/" + digest + + request, err := http.NewRequestWithContext(ctx, http.MethodDelete, url, nil) + if err != nil { + return err + } + client := s.HTTP + if client == nil { + client = &http.Client{Timeout: 30 * time.Second} + } + response, err := client.Do(request) + if err != nil { + return err + } + defer response.Body.Close() + switch response.StatusCode { + case http.StatusAccepted, http.StatusOK, http.StatusNoContent: + return nil + case http.StatusNotFound: + return Gone + case http.StatusMethodNotAllowed: + // The registry was started without deletion enabled. Said plainly, because the remedy is + // a setting on the store's module and not anything about this artifact. + return fmt.Errorf( + "the artifact store refuses deletion: its server was started without it enabled "+ + "(REGISTRY_STORAGE_DELETE_ENABLED), so nothing can be collected until the store "+ + "module is applied again (novox/hq ADR 0189). Asking about %s", reference) + default: + return fmt.Errorf("the artifact store answered %s for %s", response.Status, reference) + } +} + +// split reads a recorded path into the repository, which endpoint names the thing, and the digest. +// +// Two shapes, which are the two the mesh records: `@sha256:` is a manifest, and +// `/blobs/sha256:` is a blob. +func split(path string) (repository, kind, digest string, err error) { + if before, after, ok := strings.Cut(path, "@sha256:"); ok { + return before, "manifests", "sha256:" + after, nil + } + if before, after, ok := strings.Cut(path, "/blobs/sha256:"); ok { + return before, "blobs", "sha256:" + after, nil + } + return "", "", "", fmt.Errorf("%q names nothing the store holds by digest", path) +} diff --git a/internal/artifacts/store_test.go b/internal/artifacts/store_test.go new file mode 100644 index 0000000..3b433c8 --- /dev/null +++ b/internal/artifacts/store_test.go @@ -0,0 +1,94 @@ +package artifacts + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// Asking the store to let go of what the mesh no longer keeps (novox/hq ADR 0189, issue 108). +// +// A fake store records what it was asked to delete, so what is asserted is the mesh's decision +// and the shape of the request — not the registry's behaviour, which is the registry's to test. + +func fakeStore(t *testing.T, answer int) (Store, *[]string) { + t.Helper() + var asked []string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodDelete { + t.Errorf("the store was asked %s %s; collecting is a delete", r.Method, r.URL.Path) + } + asked = append(asked, r.URL.Path) + w.WriteHeader(answer) + })) + t.Cleanup(server.Close) + return Store{Address: strings.TrimPrefix(server.URL, "http://")}, &asked +} + +func TestAnImageAndAnArchiveAreAskedForAtTheirOwnEndpoints(t *testing.T) { + // The two shapes the mesh records: a manifest by digest, and a blob by digest. They are + // different endpoints, and asking at the wrong one answers 404 — which this would then + // record as collected, leaving the bytes on disk for ever while the record says otherwise. + store, asked := fakeStore(t, http.StatusAccepted) + ctx := context.Background() + + image := catalogue.ArtifactStoreScheme + "web/app@sha256:abc123" + archive := catalogue.ArtifactStoreScheme + "web/config/blobs/sha256:def456" + if err := store.LetGo(ctx, image); err != nil { + t.Fatal(err) + } + if err := store.LetGo(ctx, archive); err != nil { + t.Fatal(err) + } + want := []string{"/v2/web/app/manifests/sha256:abc123", "/v2/web/config/blobs/sha256:def456"} + if len(*asked) != 2 || (*asked)[0] != want[0] || (*asked)[1] != want[1] { + t.Fatalf("the store was asked %v; want %v", *asked, want) + } +} + +func TestAStoreThatDoesNotHaveItAnswersGone(t *testing.T) { + // The outcome wanted, already true. Told apart from success only so the sweep can say which + // happened; both are recorded, because retrying for ever is the thing to avoid. + store, _ := fakeStore(t, http.StatusNotFound) + err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/app@sha256:abc123") + if !errors.Is(err, Gone) { + t.Fatalf("a store that does not hold it answered %v, want Gone", err) + } +} + +func TestAStoreWithDeletionOffSaysSoAndNamesTheRemedy(t *testing.T) { + // The registry answers 405 when it was started without deletion enabled. The remedy is a + // setting on the store's module, and saying "405" would send somebody to the wrong place. + store, _ := fakeStore(t, http.StatusMethodNotAllowed) + err := store.LetGo(context.Background(), catalogue.ArtifactStoreScheme+"web/app@sha256:abc123") + if err == nil { + t.Fatal("a store that refuses deletion was read as success") + } + if !strings.Contains(err.Error(), "REGISTRY_STORAGE_DELETE_ENABLED") { + t.Fatalf("the refusal does not name the remedy: %v", err) + } +} + +func TestAReferenceThatIsNotTheMeshsOwnIsNeverAsked(t *testing.T) { + // The whole safety of the sweep is that it names only what the mesh recorded putting there. + // A reference of another shape — a vendor's image, a package version — is refused rather + // than composed into a delete somewhere that is not the mesh's store. + store, asked := fakeStore(t, http.StatusAccepted) + for _, reference := range []string{ + "docker.io/library/registry@sha256:abc123", + "registry@sha256:abc123", + "1.4.2", + } { + if err := store.LetGo(context.Background(), reference); err == nil { + t.Errorf("%s was asked about; it is not a reference into the mesh's store", reference) + } + } + if len(*asked) != 0 { + t.Fatalf("the store was asked about %v", *asked) + } +} diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index 8380cd8..dc34917 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -1511,6 +1511,13 @@ func ParseManifest(raw []byte) (Manifest, error) { } } } + // **A scheduled step may hold this module's own containers still while it runs** + // (novox/hq ADR 0189). What the host judges is the declaration it receives — whether each + // id is a container placed on that machine; what belongs here is what only the definition + // shows: that the ids are this module's, that they are containers, and that the step is + // scheduled. A module naming a neighbour's container would be a module that can stop the + // mesh, and the manifest is where that is visible. + problems = append(problems, whileStoppedProblems(m, r, hasSchedule(r))...) } for name, own := range m.OwnSecrets { if !placedOrAbsolute(own.Path) { @@ -2042,3 +2049,71 @@ func (o OwnSecrets) Paths() map[string]string { // InstancesInterchangeable is the one value of a definition's `instances`: the module is the same // on every machine, so any instance may answer for the module. const InstancesInterchangeable = "interchangeable" + +// WhileStopped is the resource key naming the containers a scheduled step holds still while it +// runs (novox/hq ADR 0189). Carried to the host unchanged, like `schedule`. +const WhileStopped = "while-stopped" + +// hasSchedule is whether a resource declares a cadence, as a string. +func hasSchedule(r map[string]any) bool { + s, _ := r["schedule"].(string) + return s != "" +} + +// whileStoppedProblems judges one container's maintenance window against its own definition +// (novox/hq ADR 0189). +// +// Three things the manifest is the only place to see: that the step is scheduled (a one-time +// offline job says *before* rather than *instead of* — at apply the host already has a window, +// because the declaration is applied in order and a run-once step gates what follows); that every +// id it names is **this module's own** container; and that it does not name itself. +// +// The host checks the fourth — that the container is actually placed on that machine — because +// that is a fact about the declaration and not about the definition. +func whileStoppedProblems(m Manifest, r map[string]any, scheduled bool) []string { + raw, present := r[WhileStopped] + if !present { + return nil + } + ids, ok := raw.([]any) + if !ok { + return []string{fmt.Sprintf( + "%s declares %s on %v as a %T; it is a list of this module's container ids", + m.Module, WhileStopped, r["id"], raw)} + } + var problems []string + if len(ids) > 0 && !scheduled { + problems = append(problems, fmt.Sprintf( + "%s declares %s on %v, which has no schedule. A maintenance window is for a recurring "+ + "step: at apply the mesh already has one, because a run-once step gates what is "+ + "declared after it (novox/hq ADR 0189)", m.Module, WhileStopped, r["id"])) + } + containers := map[string]bool{} + for _, own := range m.Resources { + if fmt.Sprint(own["type"]) == "container" { + containers[fmt.Sprint(own["id"])] = true + } + } + for _, each := range ids { + id, ok := each.(string) + if !ok { + problems = append(problems, fmt.Sprintf( + "%s declares %s on %v naming a %T; each entry is a container's id", + m.Module, WhileStopped, r["id"], each)) + continue + } + if id == fmt.Sprint(r["id"]) { + problems = append(problems, fmt.Sprintf( + "%s declares %s on %v naming itself", m.Module, WhileStopped, r["id"])) + continue + } + if !containers[id] { + problems = append(problems, fmt.Sprintf( + "%s declares %s on %v naming %q, which is not a container this module declares. "+ + "A step may hold still its own module's containers and nobody else's — one "+ + "that could quiesce a neighbour could stop the mesh", + m.Module, WhileStopped, r["id"], id)) + } + } + return problems +} diff --git a/internal/catalogue/while_stopped_test.go b/internal/catalogue/while_stopped_test.go new file mode 100644 index 0000000..39442a8 --- /dev/null +++ b/internal/catalogue/while_stopped_test.go @@ -0,0 +1,89 @@ +package catalogue + +import ( + "encoding/json" + "strings" + "testing" +) + +// A scheduled step may hold its module's own containers still while it runs (novox/hq ADR 0189). +// +// The host judges what it receives — whether each id is a container on that machine. What the +// definition is the only place to see is judged here, near whoever wrote it. + +func aStoreManifest(step map[string]any) []byte { + m := map[string]any{ + "module": "distribution", "version": "1", + "resources": []any{ + map[string]any{"id": "store", "type": "container", "name": "mesh-registry", + "image": "registry@sha256:" + strings.Repeat("a", 64)}, + step, + }, + } + raw, _ := json.Marshal(m) + return raw +} + +func TestAMaintenanceWindowOnItsOwnModulesContainerIsAccepted(t *testing.T) { + raw := aStoreManifest(map[string]any{ + "id": "collect", "type": "container", "name": "mesh-registry-collect", + "image": "registry@sha256:" + strings.Repeat("a", 64), + "schedule": "30 3 * * *", "while-stopped": []any{"store"}, + }) + if _, err := ParseManifest(raw); err != nil { + t.Fatalf("a step holding its own module's container still was refused: %v", err) + } +} + +func TestAMaintenanceWindowIsRefusedWhereTheDefinitionShowsItCannotMean(t *testing.T) { + for _, c := range []struct { + name string + step map[string]any + says string + }{ + { + "on a step with no schedule", + map[string]any{"id": "collect", "type": "container", "name": "c", + "image": "registry@sha256:" + strings.Repeat("a", 64), + "while-stopped": []any{"store"}}, + "gates what is declared after it", + }, + { + "on a run-once step, which already has order", + map[string]any{"id": "collect", "type": "container", "name": "c", + "image": "registry@sha256:" + strings.Repeat("a", 64), + "run-once": true, "while-stopped": []any{"store"}}, + "A maintenance window is for a recurring step", + }, + { + "naming a container this module does not declare", + map[string]any{"id": "collect", "type": "container", "name": "c", + "image": "registry@sha256:" + strings.Repeat("a", 64), + "schedule": "30 3 * * *", "while-stopped": []any{"the-broker"}}, + "could quiesce a neighbour could stop the mesh", + }, + { + "naming itself", + map[string]any{"id": "collect", "type": "container", "name": "c", + "image": "registry@sha256:" + strings.Repeat("a", 64), + "schedule": "30 3 * * *", "while-stopped": []any{"collect"}}, + "naming itself", + }, + { + "written as something that is not a list", + map[string]any{"id": "collect", "type": "container", "name": "c", + "image": "registry@sha256:" + strings.Repeat("a", 64), + "schedule": "30 3 * * *", "while-stopped": "store"}, + "a list of this module's container ids", + }, + } { + _, err := ParseManifest(aStoreManifest(c.step)) + if err == nil { + t.Errorf("%s was accepted", c.name) + continue + } + if !strings.Contains(err.Error(), c.says) { + t.Errorf("%s: the refusal does not say %q:\n%v", c.name, c.says, err) + } + } +} diff --git a/internal/inventory/collection.go b/internal/inventory/collection.go new file mode 100644 index 0000000..4b5bf67 --- /dev/null +++ b/internal/inventory/collection.go @@ -0,0 +1,242 @@ +package inventory + +import ( + "context" + "encoding/json" + "strings" +) + +// 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 { + if a.Reference == "" || keep[a.Reference] || collected[a.Reference] || seen[a.Reference] { + continue + } + seen[a.Reference] = true + out = append(out, a.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 a.Reference != "" { + keep[a.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 { + if a.Reference == "" || seen[a.Reference] { + continue + } + seen[a.Reference] = true + out = append(out, a.Reference) + } + } + return out, rows.Err() +} + +// 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 +} diff --git a/internal/inventory/collection_test.go b/internal/inventory/collection_test.go new file mode 100644 index 0000000..b28c3ec --- /dev/null +++ b/internal/inventory/collection_test.go @@ -0,0 +1,147 @@ +package inventory + +import ( + "context" + "fmt" + "testing" + + "github.com/novox/mesh-controller/internal/catalogue" +) + +// What the store keeps, and what it may let go (novox/hq ADR 0189, issue 108). +// +// The store has collected nothing since it was raised, and the registry's own answer — collect +// what no tag names — would delete images machines are running, because the mesh pushes under one +// moving tag and pins by digest. So the rule is the mesh's, read from its own records, and these +// are the three reasons an artifact stays and the one reason it goes. + +// ref is an artifact reference as the mesh records one. +func ref(module, artifact string, n int) string { + return fmt.Sprintf("%s%s/%s@sha256:%064x", catalogue.ArtifactStoreScheme, module, artifact, n) +} + +// built records one successful build of a module publishing one image. +func built(t *testing.T, inv *Inventory, id, module string, n int) string { + t.Helper() + reference := ref(module, "app", n) + b := aBuild(id, module, "") + b.Made = []Artifact{{Name: "app", Kind: "image", Reference: reference}} + if err := inv.RecordBuild(context.Background(), b); err != nil { + t.Fatal(err) + } + return reference +} + +func TestTheStoreKeepsTheRecentBuildsAndLetsGoOfTheRest(t *testing.T) { + inv := fresh(t) + ctx := context.Background() + + // Eight builds of one module, oldest first. Five are kept — the newest, and the four a + // release that turns out wrong can be taken back to. + var made []string + for i := 1; i <= 8; i++ { + made = append(made, built(t, inv, fmt.Sprintf("b%02d", i), "web", i)) + } + + go_, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + want := made[:3] // the three oldest + if len(go_) != len(want) { + t.Fatalf("offered %v to collect; want the %d oldest of %d", go_, len(want), len(made)) + } + for i := range want { + if go_[i] != want[i] { + t.Fatalf("offered %v; want %v — and in that order, so a failed sweep is safe to run again", + go_, want) + } + } +} + +func TestADefinitionNamingAnArtifactKeepsItHoweverOldItIs(t *testing.T) { + // The floor: no age limit. A module recorded at an older commit still names what the mesh + // would hand a machine now, and that is what must not be collected out from under it. + 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)) + } + oldest := made[0] + + // A definition the mesh holds, whose container runs that oldest image. + m := catalogue.Manifest{Module: "web", Version: "1", Resources: []map[string]any{{ + "id": "app", "type": "container", "name": "web", "image": oldest, + }}} + if err := inv.RegisterModule(ctx, m, Source{Repository: "https://forge.invalid/web.git"}); err != nil { + t.Fatal(err) + } + + go_, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + for _, reference := range go_ { + if reference == oldest { + t.Fatalf("the mesh offered to collect %s, which a definition it holds names", oldest) + } + } + if len(go_) != 2 { + t.Fatalf("offered %v; want the two oldest that nothing names", go_) + } +} + +func TestWhatHasBeenCollectedIsNotOfferedAgain(t *testing.T) { + // Without this the sweep reissues a delete for every artifact it has ever collected, every + // time it runs, for ever — a number of requests that grows with the mesh's whole history. + inv := fresh(t) + ctx := context.Background() + for i := 1; i <= 7; i++ { + built(t, inv, fmt.Sprintf("b%02d", i), "web", i) + } + first, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + if len(first) != 2 { + t.Fatalf("offered %v, want two", first) + } + if err := inv.MarkCollected(ctx, first); err != nil { + t.Fatal(err) + } + again, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + if len(again) != 0 { + t.Fatalf("offered %v again after collecting it", again) + } +} + +func TestAFailedBuildNamesNothingToCollectAndEachModuleIsCountedOnItsOwn(t *testing.T) { + inv := fresh(t) + ctx := context.Background() + + // A failed build published nothing, so it is neither kept nor collected — and it must not + // count against the module's five. + for i := 1; i <= 6; i++ { + built(t, inv, fmt.Sprintf("w%02d", i), "web", i) + } + if err := inv.RecordBuild(ctx, aBuild("w99", "web", "the recipe would not build")); err != nil { + t.Fatal(err) + } + // And a second module with three builds keeps all three: five each, not five between them. + for i := 1; i <= 3; i++ { + built(t, inv, fmt.Sprintf("d%02d", i), "db", 100+i) + } + + go_, err := inv.ToCollect(ctx) + if err != nil { + t.Fatal(err) + } + if len(go_) != 1 || go_[0] != ref("web", "app", 1) { + t.Fatalf("offered %v; want only web's oldest — db's three are all within its five", go_) + } +} diff --git a/internal/inventory/migrations/0055-what-the-store-no-longer-keeps.sql b/internal/inventory/migrations/0055-what-the-store-no-longer-keeps.sql new file mode 100644 index 0000000..28f16d7 --- /dev/null +++ b/internal/inventory/migrations/0055-what-the-store-no-longer-keeps.sql @@ -0,0 +1,23 @@ +-- What the artifact store no longer keeps (novox/hq ADR 0189, issue 108). +-- +-- The mesh removes from its store only what it put there and can account for: every digest it +-- could remove is already in a build record, so the sweep reads its own records rather than +-- enumerating the store. What it does not get from those records is whether it has already +-- removed something -- `build.made` says what that build published, for ever, which is history +-- and not an index of what is on disk. +-- +-- Without this the sweep would reissue a delete for every artifact it has ever collected, every +-- time it runs, and each one would answer 404 -- a number of requests that grows with the mesh's +-- whole history and never shrinks. +-- +-- Keyed by the reference as the mesh records it (`artifact-store:///@sha256:…`), +-- because that is the identity the record uses everywhere else. Not a foreign key to build: two +-- builds can publish the same digest (the same source built twice produces the same bytes), and +-- what is collected is the artifact, not the attempt that made it. +create table artifact_collected ( + reference text primary key, + + -- When the store answered. Kept so a reader of an old build record can tell "this artifact is + -- gone" from "this artifact was never there", which are different kinds of surprise. + at timestamptz not null default now() +);