The store keeps what the records name (hq ADR 0189)
The mesh names what may go from its own build records — a digest it did not record making is never named, which is what keeps the sweep away from the images genesis pushed. An artifact stays because a definition the mesh holds names it, or because it belongs to one of the five most recent successful builds of its module. internal/artifacts asks the store to let go of one; internal/inventory decides and remembers (migration 0055); the sweep runs after a build the mesh recorded, which is when both the bytes and the keep set moved. Never fatal to a build. And the manifest side of while-stopped, refused from the definition alone: no schedule, run-once, a container the module does not declare, itself.
This commit is contained in:
@@ -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://<module>/<artifact>@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: `<repository>@sha256:<hex>` is a manifest, and
|
||||
// `<repository>/blobs/sha256:<hex>` 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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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:<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
|
||||
}
|
||||
@@ -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_)
|
||||
}
|
||||
}
|
||||
@@ -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://<module>/<artifact>@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()
|
||||
);
|
||||
Reference in New Issue
Block a user