diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 26684d6..5ec6bca 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -558,6 +558,10 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu } 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..5464800 --- /dev/null +++ b/cmd/mesh-controller/collect.go @@ -0,0 +1,108 @@ +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. +// +// 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 + } + + // **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. Two bounds, and 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. + within, stop := context.WithTimeout(ctx, sweepBudget) + defer stop() + store := artifacts.Store{Address: address} + + var done []string + var left int + for i, reference := range references { + if i >= mostPerSweep || within.Err() != nil { + left = len(references) - i + break + } + 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. + done = append(done, reference) + continue + } + // **Stopped at the first refusal, 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. + fmt.Fprintf(os.Stderr, "the artifact store kept %s, so nothing more was asked of it: %v\n", + reference, err) + left = len(references) - i + break + } + + if len(done) > 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, done); err != nil { + 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 left > 0 { + fmt.Fprintf(os.Stderr, "%d more to collect; the next build asks again\n", left) + } +} + +// 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 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/bound_into_files.go b/internal/catalogue/bound_into_files.go index cf4ef56..b9729f3 100644 --- a/internal/catalogue/bound_into_files.go +++ b/internal/catalogue/bound_into_files.go @@ -46,7 +46,7 @@ func boundUsed(content string) [][2]string { // Three facts the mesh states about any provision, plus whatever the provider said it serves. A // module may not reach a binding it does not have — the same boundary as a secret, for the same // reason. -func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]string { +func knownFor(m Manifest, needs []Needed, node string) (map[string]map[string]string, error) { out := map[string]map[string]string{} for _, want := range m.Wants() { for i := range needs { @@ -54,12 +54,20 @@ func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]str if n.Name != want || n.For != m.Module { continue } + as := ConsumerIdentity(node, IdentitySource(m.Slug, m.Module)) values := map[string]string{ "at": n.At, "from": n.From, - "as": ConsumerIdentity(node, IdentitySource(m.Slug, m.Module)), + "as": as, } - for key, value := range n.Serves { + // What the provider derives for this consumer rather than for all of them + // (novox/hq ADR 0201). Filled here, the one place a provision and the module + // requiring it are both in hand. + served, err := ServedTo(n.Serves, as) + if err != nil { + return nil, fmt.Errorf("%s requires %s: %w", m.Module, want, err) + } + for key, value := range served { // The provider's own vocabulary. Rendered plainly: a port is 5432, not 5432.000000, // which is what a float would write and what a connection string would refuse. values[key] = plainly(value) @@ -67,7 +75,7 @@ func knownFor(m Manifest, needs []Needed, node string) map[string]map[string]str out[want] = values } } - return out + return out, nil } // withOwnNames adds a module's own composed names to what it may name from one binding: diff --git a/internal/catalogue/consumer_into_serves.go b/internal/catalogue/consumer_into_serves.go new file mode 100644 index 0000000..897411d --- /dev/null +++ b/internal/catalogue/consumer_into_serves.go @@ -0,0 +1,311 @@ +package catalogue + +import ( + "fmt" + "regexp" + "sort" + "strings" +) + +// What a provider derives for one consumer, said once in the provider's definition and delivered +// to both ends (novox/hq ADR 0201, issue 124). +// +// A `serves` block is otherwise literal: the same values for every consumer. Where the provider +// *names the resource* — a bucket, a database, a vhost — the name is derived from who is asking, +// and before this the mesh had no channel for it. The provider recomputed it in its own code and +// every consumer transcribed it into its own definition by hand, which is a copy of somebody +// else's rule kept in agreement by nobody. One of three transcriptions was wrong for months. +// +// **The mesh learns no protocol here; it spells its own name in an alphabet it already knows.** +// The only fact a served value may name is the identity the mesh itself minted for the consumer, +// in one of two alphabets: as it was minted, and as a DNS label. Everything a provider wants +// around it — a prefix, a suffix, a separator — it writes around the placeholder, because a +// served value is a string. + +// consumerFact is `${consumer:}` or `${consumer::}`. +var consumerFact = regexp.MustCompile(`\$\{consumer:([a-z][a-z0-9-]*)(?::([a-z][a-z0-9-]*))?\}`) + +// consumerFacts are what a served value may name about the consumer it is being derived for. +// One entry, deliberately: the identity is the one thing about a consumer the mesh itself chose, +// so it is the one thing the mesh can hand to a provider without either end guessing. +var consumerFacts = []string{"as"} + +// consumerAlphabets are the ways the mesh will write that identity. `dns` is the mesh's own +// identifier with its separator written `-` instead of `_` — the whole of the difference between +// the alphabet the mesh mints in and the one buckets, vhosts and hostnames accept. +var consumerAlphabets = []string{"dns"} + +// ServedTo fills a provider's served values for one consumer. +// +// `as` is the identity the mesh minted for that consumer — the same string it is told to present +// as a login. Values with no placeholder are returned exactly as they were, and a block with no +// placeholder at all is returned unchanged, so this costs nothing for the providers that derive +// nothing. +// +// Only strings carry placeholders. A number, a boolean or a nested object is a value the provider +// stated outright, and is left alone. +func ServedTo(serves map[string]any, as string) (map[string]any, error) { + if len(serves) == 0 { + return serves, nil + } + var out map[string]any + for _, key := range sortedAnyKeys(serves) { + text, ok := serves[key].(string) + if !ok || !strings.Contains(text, "${consumer:") { + continue + } + filled, err := consumerInto(text, as) + if err != nil { + return nil, fmt.Errorf("the value served as %q: %w", key, err) + } + if out == nil { + // Copied only once something actually changes: the caller's map is the manifest's, + // and a provider that derives nothing must not have it rewritten underneath it. + out = make(map[string]any, len(serves)) + for k, v := range serves { + out[k] = v + } + } + out[key] = filled + } + if out == nil { + return serves, nil + } + return out, nil +} + +// consumerInto replaces every `${consumer:…}` in one value. +// +// **A fact or an alphabet the mesh does not have is refused, not left standing.** Written through, +// the literal `${consumer:as}` would reach a configuration file and be read as a bucket name, +// failing somewhere that names neither the module nor the mesh — the same reasoning `${bound:…}` +// is refused by (boundInto). +func consumerInto(value, as string) (string, error) { + var failed error + out := consumerFact.ReplaceAllStringFunc(value, func(match string) string { + parts := consumerFact.FindStringSubmatch(match) + fact, alphabet := parts[1], parts[2] + if fact != "as" { + if failed == nil { + failed = fmt.Errorf( + "says %s, and the mesh states %s about a consumer", match, orNothing(consumerFacts)) + } + return match + } + switch alphabet { + case "": + return as + case "dns": + return asDNSLabel(as) + default: + if failed == nil { + failed = fmt.Errorf( + "says %s, and the mesh writes an identity as %s", match, orNothing(consumerAlphabets)) + } + return match + } + }) + if failed != nil { + return "", failed + } + return out, nil +} + +// asDNSLabel writes a minted identity as a DNS label. +// +// The mesh's identities are already lower-case letters, digits and `_` (ConsumerIdentity), and +// already short enough for the tightest backend they reach (CheckIdentity, twenty characters). So +// this is the separator and nothing else — no lower-casing of what is already lower case, no +// truncation to a limit the identity is already inside, no padding of a name that is already long +// enough. Each of those would be the mesh guessing at a rule it has not been given. +func asDNSLabel(as string) string { + return strings.ReplaceAll(as, "_", "-") +} + +// CheckServes refuses a `serves` block that names a consumer fact or an alphabet the mesh does not +// have, when the definition is parsed rather than when a consumer is resolved. +// +// A provision nobody consumes yet still has its rule read: a definition that would be refused the +// first time somebody required it is a definition that is wrong now. +func CheckServes(m Manifest) []string { + var problems []string + for _, provision := range sortedServes(m.Serves) { + for _, key := range sortedAnyKeys(m.Serves[provision]) { + text, ok := m.Serves[provision][key].(string) + if !ok { + continue + } + // A probe identity, because what is checked is the shape of the statement and not + // what any consumer is called. + if _, err := consumerInto(text, "mesh_node_module"); err != nil { + problems = append(problems, fmt.Sprintf( + "%s serves %s, and the value it serves as %q %s", m.Module, provision, key, err)) + } + } + } + return problems +} + +func sortedServes(serves map[string]map[string]any) []string { + out := make([]string, 0, len(serves)) + for k := range serves { + out = append(out, k) + } + sort.Strings(out) + return out +} + +func sortedAnyKeys(values map[string]any) []string { + out := make([]string, 0, len(values)) + for k := range values { + out = append(out, k) + } + sort.Strings(out) + return out +} + +// derivedFor is what the provider on this machine derives for one consumer of one provision +// (novox/hq ADR 0201). +// +// Settled first, then derived: an operator may set a prefix on what the provider serves and the +// mesh still fills the consumer's half of it ([ADR 0174]). Only the keys that actually name the +// consumer are returned — the rest of a `serves` block is the same for every consumer and is +// already in the provider's own definition, so repeating it here would be a second copy to go +// stale. +// +// The first module in the resolved order that says it serves the provision answers, which is the +// choice servedOnThisMachine makes for the consumer's half. Nothing serving it on this machine is +// not an error: a contribution can reach a machine whose provider is a record or an adapter, and +// then there is nothing derived to tell. +func (r Resolution) derivedFor(provision, as, consumer, local string, settings SettingsBy) (map[string]any, error) { + for _, m := range r.Modules { + serves, said := m.Serves[provision] + if !said { + continue + } + var names map[string]any + for key, value := range serves { + if text, ok := value.(string); ok && strings.Contains(text, "${consumer:") { + if names == nil { + names = map[string]any{} + } + names[key] = value + } + } + if names == nil { + return nil, nil + } + // **A consumer that keeps several holders of this provision is refused** — this is issue + // 124's own failure one case to the side, and it would be just as quiet. + // + // Each holder gets its own login, `…_` (ADR 0094), and a provider derives from the + // login, so it would make one resource per holder. The consumer's side has no such + // dimension: one binding file per provision, one `${bound::}`, both + // derived from the un-suffixed identity. So the provider would create the holder's + // resource and the consumer would be configured against a name nothing made — it would + // authenticate successfully and be refused on every object, which reads like a credential + // fault and is not one. + // + // Lifting this means giving the consumer's side a local dimension. That is a decision, + // not an omission, and until it is taken the mesh says so rather than guessing. + if local != "" { + return nil, fmt.Errorf( + "%s keeps several holders of %s (this one is %q), and %s derives %s for each "+ + "consumer from the login the mesh minted. Each holder has its own login, and a "+ + "consumer is told one value per requirement — so the two ends would name "+ + "different things and nothing would compare them (novox/hq ADR 0201)", + consumer, local, provision, m.Module, orNothing(sortedAnyKeys(names))) + } + settled, err := Settle(names, settings[m.Module]) + if err != nil { + return nil, fmt.Errorf("%s serving %s: %w", m.Module, provision, err) + } + derived, err := ServedTo(settled, as) + if err != nil { + return nil, fmt.Errorf("%s serving %s to %s: %w", m.Module, provision, as, err) + } + return derived, nil + } + return nil, nil +} + +// notTranscribed refuses a consumer's file that writes out the value its provider derives for it, +// instead of asking for it (novox/hq ADR 0201, issue 124). +// +// **What would have caught the one wrong instance.** The object store's three consumers each wrote +// their bucket into their own configuration by hand. One of them named a predecessor's bucket, and +// nothing compared it to what the provider would actually create: the module would have +// authenticated successfully and been refused on every object, which reads like a credential fault +// and is not one. It looked authoritative for months. +// +// The test is exact and costs one string search: a definition whose file already contains the +// value the mesh is about to derive for it has written down somebody else's rule. It cannot be a +// coincidence — a derived value carries the identity the mesh minted for this very consumer on +// this very machine, which nothing else would spell out — and it cannot be checked afterwards, +// because after substitution every consumer's file contains it legitimately. +// +// Only values that actually name the consumer are judged. A provider that serves a constant under +// the same key serves the same constant to everyone, and a consumer repeating it is redundant +// rather than wrong. +func notTranscribed(resource map[string]any, known map[string]map[string]string, module string) error { + if fmt.Sprint(resource["type"]) != "file" { + return nil + } + content, ok := resource["content"].(string) + if !ok || content == "" { + return nil + } + for _, provision := range sortedKnown(known) { + values := known[provision] + identity := values["as"] + if identity == "" { + continue + } + for _, key := range sortedStringKeys(values) { + if key == "as" { + // The login is not derived from itself, and a consumer that must present it in a + // connection string legitimately has it from `${bound:…}` — which is what it will + // be after substitution, so this would judge the substitution, not the module. + continue + } + value := values[key] + if value == "" || !namesTheConsumer(value, identity) { + continue + } + if !strings.Contains(content, value) { + continue + } + return fmt.Errorf( + "%s writes %q into %v, and that is exactly what %s derives for it — a definition "+ + "keeping its own copy of somebody else's naming rule is one that can disagree "+ + "with it, silently. Say ${bound:%s:%s} and be told", + module, value, resource["id"], provision, provision, key) + } + } + return nil +} + +// namesTheConsumer is whether a derived value was built from this consumer's identity — in the +// alphabet it was minted in, or as a DNS label. A value that does not contain it was not derived +// from it, whatever else it may be. +func namesTheConsumer(value, identity string) bool { + return strings.Contains(value, identity) || strings.Contains(value, asDNSLabel(identity)) +} + +func sortedKnown(known map[string]map[string]string) []string { + out := make([]string, 0, len(known)) + for k := range known { + out = append(out, k) + } + sort.Strings(out) + return out +} + +func sortedStringKeys(values map[string]string) []string { + out := make([]string, 0, len(values)) + for k := range values { + out = append(out, k) + } + sort.Strings(out) + return out +} diff --git a/internal/catalogue/declaration.go b/internal/catalogue/declaration.go index d69eb7a..b2f651b 100644 --- a/internal/catalogue/declaration.go +++ b/internal/catalogue/declaration.go @@ -645,7 +645,16 @@ func (r Resolution) compose(with Rendering, owner map[string]string, if err != nil { return nil, err } - file, err := boundFile(*found, m.Binds[to], ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)), own) + as := ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module)) + // What the provider derives for THIS consumer, filled here where the consumer is + // known (novox/hq ADR 0201). The same fill knownFor does below, so the binding file + // and the module's `${bound:…}` substitutions cannot say different things. + told := *found + told.Serves, err = ServedTo(told.Serves, as) + if err != nil { + return nil, fmt.Errorf("%s is told about %s: %w", m.Module, to, err) + } + file, err := boundFile(told, m.Binds[to], as, own) if err != nil { return nil, err } @@ -719,7 +728,10 @@ func (r Resolution) compose(with Rendering, owner map[string]string, return nil, err } // And what its bindings say, for the half of a connection that is not secret. - known := knownFor(m, r.Needs, r.Node) + known, err := knownFor(m, r.Needs, r.Node) + if err != nil { + return nil, err + } // A requirement answered on this same machine is not in r.Needs — its binding file is // written from `here` (above) — and so `${bound:…}` could not name it, though the file // beside it said the same facts. Filled from the same answer, so the two cannot disagree. @@ -736,7 +748,11 @@ func (r Resolution) compose(with Rendering, owner map[string]string, } local := *answered local.For = m.Module - for provision, values := range knownFor(m, []Needed{local}, r.Node) { + here, err := knownFor(m, []Needed{local}, r.Node) + if err != nil { + return nil, err + } + for provision, values := range here { known[provision] = values } } @@ -761,6 +777,17 @@ func (r Resolution) compose(with Rendering, owner map[string]string, // And the machine underneath, which no binding of its own can tell it. thisMachine := machineFacts(r, with.Names, with.MeshRange) + // **A definition that already holds the answer transcribed it** (novox/hq ADR 0201). + // Judged over what the module itself declares, and before anything is substituted: the + // mesh's own generated files — the binding, the contributions — legitimately carry the + // derived value, and after substitution so does every consumer's file, so this is the one + // moment the two can be told apart. + for _, own := range m.Resources { + if err := notTranscribed(own, known, m.Module); err != nil { + return nil, err + } + } + // Which of this module's files carry a secret, for the rule that a container may not read // one of them as its environment without saying so (ADR 0086, issue 041). secretFiles := secretFilesOf(resources) @@ -1175,6 +1202,19 @@ type Contribution struct { // requirement's name — everything providing `reverse-proxy` understands the same shape, which // is what makes swapping one for another cost nothing. Values map[string]any `json:"values"` + // Derived is what this provider's own definition said it derives for this consumer, already + // derived (novox/hq ADR 0201). + // + // **The provider is told, rather than recomputing it.** A served value may name the consumer's + // identity — a bucket named for who is asking, a database prefixed with it — and before this + // the rule lived twice: once in the provisioner's code, once transcribed into every consumer's + // definition. The mesh fills the provider's own statement here and delivers the same filled + // value to the consumer, so the two cannot disagree: there is no second computation to + // disagree with. + // + // Only the keys that are per-consumer. The rest of what the provider serves is the same for + // everyone and is in its own definition, where it already is. + Derived map[string]any `json:"derived,omitempty"` } // grantPath is where one consumer's sealed credential lands on the providing machine. @@ -1266,12 +1306,17 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant, // told about it and withdraws the login on its next pass. continue } + as := holderAs(ConsumerIdentity(g.Consumer, IdentitySource(g.Slug, g.From)), g.Local) + derived, err := r.derivedFor(g.Provision, as, g.From, g.Local, settings) + if err != nil { + return nil, err + } out[g.Provision] = append(out[g.Provision], Contribution{ - From: g.From, Node: g.Consumer, At: g.At, Values: g.Values, + From: g.From, Node: g.Consumer, At: g.At, Values: g.Values, Derived: derived, // One holder per local name: the identity the consumer is known by, and the local name // after it where the module keeps several (ADR 0094). Not a login any backend checks — // a secret is not a login — so the identity limit does not apply to the suffix. - As: holderAs(ConsumerIdentity(g.Consumer, IdentitySource(g.Slug, g.From)), g.Local), + As: as, Secret: grantPath(directories[g.Provision], g.Consumer, holderAs(g.From, g.Local)), }) if granted[g.Provision] == nil { diff --git a/internal/catalogue/derived_for_consumer_test.go b/internal/catalogue/derived_for_consumer_test.go new file mode 100644 index 0000000..5d85206 --- /dev/null +++ b/internal/catalogue/derived_for_consumer_test.go @@ -0,0 +1,343 @@ +package catalogue + +import ( + "encoding/json" + "strings" + "testing" +) + +// What a provider derives for each consumer, said once and delivered to both ends +// (novox/hq ADR 0201, issue 124). +// +// The failure these are written against: the object store's provisioner derived each consumer's +// bucket from the login the mesh minted, in its own code, and the mesh had no channel to tell the +// consumer which bucket that was — so all three consumers wrote the answer into their own +// definitions by hand. Two were right. One named a predecessor's bucket and would have +// authenticated successfully and been refused on every object. Each of them also named the +// machine the module happens to run on, which a definition may not do. + +// store is an object store in the shape minio has: it serves a region and a port to everyone, and +// a bucket named for whoever is asking. +func store() Manifest { + return Manifest{ + Module: "store", Version: "1", + Provides: FromAnywhere("s3-bucket"), + Listens: []Listening{{Port: 9000, Protocol: "tcp", From: FromMesh}}, + Serves: map[string]map[string]any{"s3-bucket": { + "region": "eu-west", + "bucket": "${consumer:as:dns}", + }}, + Receives: map[string]string{"s3-bucket": "/var/lib/store/grants/mesh.json"}, + Grants: map[string]string{"s3-bucket": "/var/lib/store/grants"}, + Resources: []map[string]any{{ + "id": "server", "type": "container", "name": "store", "ports": []any{"9000"}, + }}, + } +} + +// files is a consumer that writes the bucket into its own configuration — which is the thing it +// could not do before, and had to transcribe. +func files() Manifest { + return Manifest{ + Module: "files", Version: "1", Slug: "files", + Requires: []string{"s3-bucket"}, + Binds: map[string]string{"s3-bucket": "/var/lib/files/store.json"}, + Secrets: map[string]string{"s3-bucket": "/var/lib/files/store.secret"}, + Resources: []map[string]any{{ + "id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600", + "content": "BUCKET=${bound:s3-bucket:bucket}\nREGION=${bound:s3-bucket:region}\n", + }}, + } +} + +// pics is a second consumer of the same provider on the same machine: two derivations, neither +// the other's. +func pics() Manifest { + return Manifest{ + Module: "pics", Version: "1", Slug: "pics", + Requires: []string{"s3-bucket"}, + Binds: map[string]string{"s3-bucket": "/var/lib/pics/store.json"}, + Secrets: map[string]string{"s3-bucket": "/var/lib/pics/store.secret"}, + Resources: []map[string]any{{ + "id": "env", "type": "file", "path": "/var/lib/pics/env", "mode": "0600", + "content": "BUCKET=${bound:s3-bucket:bucket}\n", + }}, + } +} + +// The three places the derived value lands must agree, because agreeing is the whole point: the +// consumer's own file, the binding it reads as JSON, and the provider's contributions entry. +func TestADerivedValueReachesBothEndsAndAgrees(t *testing.T) { + r, err := Resolve(shelf(store(), files()), []string{"store", "files"}, reachable(), World{}) + if err != nil { + t.Fatal(err) + } + out, err := r.Declaration(Rendering{Grants: []Grant{{ + Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Values: map[string]any{}, Sealed: "c2VhbGVk", + }}}) + if err != nil { + t.Fatal(err) + } + + // The mesh minted this identity for the consumer; the bucket is that identity as a DNS label. + // Derived here with the mesh's own function, so the test cannot agree with a wrong rule. + as := ConsumerIdentity("workstation", IdentitySource("files", "files")) + want := strings.ReplaceAll(as, "_", "-") + if want == as || !strings.Contains(as, "_") { + t.Fatalf("the mesh's identity %q has no separator to rewrite; this test proves nothing", as) + } + + env := fileNamed(out, "files.env") + if env == nil { + t.Fatalf("the consumer was given no file: %v", out) + } + if got := env["content"].(string); !strings.Contains(got, "BUCKET="+want+"\n") { + t.Errorf("the consumer's own file was not told the bucket:\n%s\nwant BUCKET=%s", got, want) + } + + binding := fileNamed(out, "files.bound-s3-bucket") + if binding == nil { + t.Fatalf("the consumer was given no binding: %v", out) + } + var said struct { + Serves map[string]any `json:"serves"` + } + if err := json.Unmarshal([]byte(binding["content"].(string)), &said); err != nil { + t.Fatal(err) + } + if said.Serves["bucket"] != want { + t.Errorf("the binding says the bucket is %q, want %q", said.Serves["bucket"], want) + } + // And what is the same for everybody is still the same for everybody. + if said.Serves["region"] != "eu-west" { + t.Errorf("the binding lost what the provider serves to all: %v", said.Serves) + } + + given := storeGrants(t, out) + if len(given) != 1 { + t.Fatalf("the provider was told about %d consumer(s): %v", len(given), given) + } + if given[0].Derived["bucket"] != want { + t.Errorf("the provider was told the bucket is %v, and the consumer was told %q — "+ + "the two ends disagree, which is the whole failure", given[0].Derived["bucket"], want) + } + // Only the per-consumer half. The region is the same for everyone and is already in the + // provider's own definition; repeating it here would be a copy to go stale. + if _, carried := given[0].Derived["region"]; carried { + t.Errorf("the provider was handed back what it already says for everyone: %v", given[0].Derived) + } +} + +// Two consumers of one provider on one machine get two buckets, and neither gets the other's. +func TestTwoConsumersOfOneProviderGetTheirOwnDerivation(t *testing.T) { + r, err := Resolve(shelf(store(), files(), pics()), + []string{"store", "files", "pics"}, reachable(), World{}) + if err != nil { + t.Fatal(err) + } + out, err := r.Declaration(Rendering{Grants: []Grant{ + {Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Values: map[string]any{}, Sealed: "c2VhbGVk"}, + {Provision: "s3-bucket", Consumer: "workstation", From: "pics", Slug: "pics", + Values: map[string]any{}, Sealed: "c2VhbGVk"}, + }}) + if err != nil { + t.Fatal(err) + } + forFiles := strings.ReplaceAll(ConsumerIdentity("workstation", IdentitySource("files", "files")), "_", "-") + forPics := strings.ReplaceAll(ConsumerIdentity("workstation", IdentitySource("pics", "pics")), "_", "-") + if forFiles == forPics { + t.Fatal("the two consumers were given the same identity; this test proves nothing") + } + if got := fileNamed(out, "files.env")["content"].(string); !strings.Contains(got, "BUCKET="+forFiles+"\n") { + t.Errorf("files was not given its own bucket:\n%s", got) + } + if got := fileNamed(out, "pics.env")["content"].(string); !strings.Contains(got, "BUCKET="+forPics+"\n") { + t.Errorf("pics was not given its own bucket:\n%s", got) + } + var buckets []any + for _, g := range storeGrants(t, out) { + buckets = append(buckets, g.Derived["bucket"]) + } + if len(buckets) != 2 || buckets[0] == buckets[1] { + t.Errorf("the provider was told %v; it must be told one bucket per consumer", buckets) + } +} + +// An operator may still set what the provider serves, and the mesh still derives the rest: the +// setting is laid on first, then the consumer's half is filled. +func TestASettingComposesWithADerivedValue(t *testing.T) { + r, err := Resolve(shelf(store(), files()), []string{"store", "files"}, reachable(), World{}) + if err != nil { + t.Fatal(err) + } + out, err := r.Declaration(Rendering{ + Settings: SettingsBy{"store": {{From: "the operator", + Values: map[string]any{"bucket": "team-${consumer:as:dns}"}}}}, + Grants: []Grant{{Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Values: map[string]any{}, Sealed: "c2VhbGVk"}}, + }) + if err != nil { + t.Fatal(err) + } + want := "team-" + strings.ReplaceAll(ConsumerIdentity("workstation", IdentitySource("files", "files")), "_", "-") + if got := fileNamed(out, "files.env")["content"].(string); !strings.Contains(got, "BUCKET="+want+"\n") { + t.Errorf("the operator's prefix did not survive the derivation:\n%s\nwant BUCKET=%s", got, want) + } + if given := storeGrants(t, out); given[0].Derived["bucket"] != want { + t.Errorf("the provider was told %v, the consumer %q", given[0].Derived["bucket"], want) + } +} + +// A fact or an alphabet the mesh does not have is refused where the definition is, not where a +// consumer happens to be resolved — and the refusal says what may be said instead. +func TestAServedValueNamingSomethingTheMeshDoesNotHaveIsRefused(t *testing.T) { + for _, c := range []struct{ value, says string }{ + {"${consumer:node}", "as"}, + {"${consumer:as:punycode}", "dns"}, + } { + m := store() + m.Serves["s3-bucket"]["bucket"] = c.value + raw, err := json.Marshal(m) + if err != nil { + t.Fatal(err) + } + _, err = ParseManifest(raw) + if err == nil { + t.Fatalf("%s was accepted", c.value) + } + if !strings.Contains(err.Error(), c.value) { + t.Errorf("the refusal of %s does not quote it: %v", c.value, err) + } + if !strings.Contains(err.Error(), c.says) { + t.Errorf("the refusal of %s does not say what may be said (%q): %v", c.value, c.says, err) + } + } +} + +// `dns` is checked against an identity the mesh actually mints, not an invented string. +func TestTheDNSAlphabetIsTheMintedIdentityWithItsSeparatorRewritten(t *testing.T) { + as := ConsumerIdentity("anchor", IdentitySource("ncloud", "nextcloud")) + if err := CheckIdentity("anchor", IdentitySource("ncloud", "nextcloud")); err != nil { + t.Fatalf("the mesh would not mint this identity at all: %v", err) + } + label := asDNSLabel(as) + if strings.Contains(label, "_") { + t.Errorf("%q is not a DNS label", label) + } + if strings.ReplaceAll(label, "-", "_") != as { + t.Errorf("%q is not %q with its separator rewritten", label, as) + } +} + +// The check that would have caught the one wrong instance: a consumer that writes the derived +// value into its own definition instead of asking for it is refused, whether it transcribed the +// right answer or a predecessor's. +func TestAConsumerThatTranscribesWhatItsProviderDerivesIsRefused(t *testing.T) { + as := ConsumerIdentity("workstation", IdentitySource("files", "files")) + transcribed := strings.ReplaceAll(as, "_", "-") + + m := files() + m.Resources = []map[string]any{{ + "id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600", + // Exactly what the provider will create — correct today, and a copy of a rule that is + // not this module's. + "content": "BUCKET=" + transcribed + "\n", + }} + r, err := Resolve(shelf(store(), m), []string{"store", "files"}, reachable(), World{}) + if err != nil { + t.Fatal(err) + } + _, err = r.Declaration(Rendering{Grants: []Grant{{ + Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Values: map[string]any{}, Sealed: "c2VhbGVk", + }}}) + if err == nil { + t.Fatal("a definition holding its own copy of the provider's naming rule was accepted") + } + if !strings.Contains(err.Error(), "${bound:s3-bucket:bucket}") { + t.Errorf("the refusal does not say what to write instead: %v", err) + } + + // And a constant the provider serves to everyone is not a transcription: repeating it is + // redundant, not wrong, and refusing it would be the mesh policing style. + m.Resources = []map[string]any{{ + "id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600", + "content": "REGION=eu-west\n", + }} + r, err = Resolve(shelf(store(), m), []string{"store", "files"}, reachable(), World{}) + if err != nil { + t.Fatal(err) + } + if _, err := r.Declaration(Rendering{Grants: []Grant{{ + Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Values: map[string]any{}, Sealed: "c2VhbGVk", + }}}); err != nil { + t.Errorf("a value the provider serves to everyone was judged a transcription: %v", err) + } +} + +func storeGrants(t *testing.T, out []map[string]any) []Contribution { + t.Helper() + for _, r := range out { + if r["path"] != "/var/lib/store/grants/mesh.json" { + continue + } + var parsed struct { + Given []Contribution `json:"given"` + } + if err := json.Unmarshal([]byte(r["content"].(string)), &parsed); err != nil { + t.Fatal(err) + } + return parsed.Given + } + t.Fatalf("the provider was given no contributions file: %v", out) + return nil +} + +// A consumer that keeps SEVERAL holders of one provision is refused, rather than told one thing +// while its provider is told another. +// +// **This is issue 124's own failure, one case to the side.** The mesh gives each holder its own +// login — `mesh_node_mod_` (ADR 0094) — and the provider derives from the login, so it +// would make one resource per holder. The consumer's side has no such dimension: there is one +// binding file per provision and one `${bound::}`, both derived from the +// un-suffixed identity. So the provider would create `…-mod-cold` and the consumer would be +// configured against `…-mod`: it would authenticate successfully and be refused on every object, +// which is exactly the fault this whole record exists to end. +// +// Refused, loudly, at the one place that can see both halves. Lifting it means giving the +// consumer's side a local dimension, which is a decision and not an omission. +func TestAConsumerWithSeveralHoldersOfADerivingProviderIsRefused(t *testing.T) { + m := files() + // Two holders of the one provision, the shape ADR 0094 gives a module that keeps several. + m.Secrets = nil + m.SecretsMany = map[string]map[string]string{"s3-bucket": { + "hot": "/var/lib/files/hot.secret", + "cold": "/var/lib/files/cold.secret", + }} + m.Resources = []map[string]any{{ + "id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600", + "content": "BUCKET=${bound:s3-bucket:bucket}\n", + }} + r, err := Resolve(shelf(store(), m), []string{"store", "files"}, reachable(), World{}) + if err != nil { + t.Fatal(err) + } + _, err = r.Declaration(Rendering{Grants: []Grant{ + {Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Local: "hot", Values: map[string]any{}, Sealed: "c2VhbGVk"}, + {Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files", + Local: "cold", Values: map[string]any{}, Sealed: "c2VhbGVk"}, + }}) + if err == nil { + t.Fatal("a consumer with several holders of a deriving provider was accepted; " + + "its two ends would have disagreed in silence") + } + for _, want := range []string{"files", "s3-bucket", "bucket"} { + if !strings.Contains(err.Error(), want) { + t.Errorf("the refusal does not name %q: %v", want, err) + } + } +} diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index 8528c9f..9dd482e 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -1428,6 +1428,10 @@ func ParseManifest(raw []byte) (Manifest, error) { "%s serves %q to whoever requires it, and does not provide it", m.Module, to)) } } + // A served value may be derived for the consumer it is served to (novox/hq ADR 0201). Read + // here, where the definition is, rather than when somebody first requires it: a rule that + // would be refused at the first consumer is wrong from the moment it is written. + problems = append(problems, CheckServes(m)...) for to, where := range m.Binds { if !placedOrAbsolute(where) { problems = append(problems, fmt.Sprintf( @@ -1622,6 +1626,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) { @@ -2153,3 +2164,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/0056-what-the-store-no-longer-keeps.sql b/internal/inventory/migrations/0056-what-the-store-no-longer-keeps.sql new file mode 100644 index 0000000..28f16d7 --- /dev/null +++ b/internal/inventory/migrations/0056-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() +);