Carry an adopted node's mode and taken modules in every declaration, from one marshaller (hq ADR 0100)

This commit is contained in:
2026-09-22 17:17:54 +02:00
parent 1c32af6a22
commit a86a6c2974
6 changed files with 357 additions and 49 deletions
+36 -19
View File
@@ -1,6 +1,7 @@
package main
import (
"bytes"
"context"
"encoding/json"
"errors"
@@ -271,10 +272,10 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
// silently produced a declaration missing them — a difference between what `plan` showed and what
// `plan --json` handed to anything reading it.
func declarationFor(ctx context.Context, open *stores, node string,
plan catalogue.Resolution, settings catalogue.SettingsBy) ([]map[string]any, error) {
plan catalogue.Resolution, settings catalogue.SettingsBy) (sendable, error) {
gens, err := generators(ctx, open)
if err != nil {
return nil, err
return sendable{}, err
}
// **Allocating, because `plan` is the send without the sending.** It is one machine, named by
// a person, who is asking what a push would do — so the port it shows and the secret it seals
@@ -316,11 +317,11 @@ const (
func declarationWith(ctx context.Context, open *stores, node string,
plan catalogue.Resolution, settings catalogue.SettingsBy,
gens map[string]catalogue.Generator, choosing Choosing) ([]map[string]any, error) {
gens map[string]catalogue.Generator, choosing Choosing) (sendable, error) {
inv := open.inventory
grants, err := grantsFor(ctx, open, node)
if err != nil {
return nil, err
return sendable{}, err
}
// Where this machine puts what each module needs reachable (novox/hq ADR 0038).
//
@@ -333,7 +334,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
if choosing == Reading {
held, err := inv.PortsFor(ctx, node)
if err != nil {
return nil, err
return sendable{}, err
}
for _, a := range held {
if already[a.Module] == nil {
@@ -357,7 +358,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
case mayAssign && choosing == Allocating:
at, err := inv.PortFor(ctx, node, m.Module, l.Port, l.Fixed)
if err != nil {
return nil, fmt.Errorf(
return sendable{}, fmt.Errorf(
"%s needs %d reachable on %s and it could not be assigned: %w",
m.Module, l.Port, node, err)
}
@@ -398,7 +399,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
}
}
if err != nil {
return nil, err
return sendable{}, err
}
if needed[m.Module] == nil {
needed[m.Module] = map[string]string{}
@@ -416,7 +417,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
}
issued, meshCA, err := certificateFor(ctx, open, node)
if err != nil {
return nil, err
return sendable{}, err
}
certificate, authority = issued, meshCA
break
@@ -429,11 +430,11 @@ func declarationWith(ctx context.Context, open *stores, node string,
// One reading of the catalogue for the three questions below that resolve the whole mesh.
shelf, err := inv.Catalogue(ctx)
if err != nil {
return nil, err
return sendable{}, err
}
private, err := onThePrivateNetwork(ctx, inv, shelf)
if err != nil {
return nil, err
return sendable{}, err
}
// And every machine's name, so a container can reach one. The same set that writes the
@@ -441,7 +442,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
// about where another machine is.
names, err := namesInTheMesh(ctx, inv, shelf)
if err != nil {
return nil, err
return sendable{}, err
}
// And every routed name → the node that serves it (novox/hq ADR 0066). Alongside the
@@ -450,7 +451,7 @@ func declarationWith(ctx context.Context, open *stores, node string,
// to serve and knows nothing about what they mean.
routes, err := routeNamesInTheMesh(ctx, open)
if err != nil {
return nil, err
return sendable{}, err
}
for name, at := range routes {
names[name] = at
@@ -479,15 +480,25 @@ func declarationWith(ctx context.Context, open *stores, node string,
continue
}
if kept, err = inv.OperatorExport(ctx); err != nil {
return nil, err
return sendable{}, err
}
break
}
return plan.Declaration(catalogue.Rendering{
composed, err := plan.Compose(catalogue.Rendering{
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
Certificate: certificate, Authority: authority, Mesh: private, Names: names,
Suffix: overlay.Suffix(), Foundation: foundation, Kept: kept})
if err != nil {
return sendable{}, err
}
// Whether this node is adopted, and what was taken on it, said in every declaration it is
// sent from this one place (novox/hq ADR 0100).
adoption, err := adoptionOf(ctx, inv, node, plan, composed)
if err != nil {
return sendable{}, err
}
return sendable{Resources: composed.Resources, Adoption: adoption}, nil
}
// routeNamesInTheMesh is every routed name and the address of the node that serves it (novox/hq
@@ -736,16 +747,21 @@ func planCommand(ctx context.Context, args []string) error {
return nil
}
if *asJSON {
resources, err := declarationFor(ctx, open, args[0], plan, settings)
declared, err := declarationFor(ctx, open, args[0], plan, settings)
if err != nil {
return err
}
body, err := json.MarshalIndent(
map[string]any{"declaration": 1, "resources": resources}, "", " ")
// The bytes a push would send, indented: one marshaller, so `plan --json` cannot show an
// envelope other than the one sent.
body, err := declared.Body()
if err != nil {
return err
}
fmt.Println(string(body))
var indented bytes.Buffer
if err := json.Indent(&indented, body, "", " "); err != nil {
return err
}
fmt.Println(indented.String())
return nil
}
@@ -769,10 +785,11 @@ func planCommand(ctx context.Context, args []string) error {
for _, n := range plan.Needs {
fmt.Printf(" needs %s from %s, for %s\n", n.Name, n.From, n.For)
}
resources, err := declarationFor(ctx, open, args[0], plan, settings)
declared, err := declarationFor(ctx, open, args[0], plan, settings)
if err != nil {
return err
}
resources := declared.Resources
for module, layers := range settings {
for _, layer := range layers {
fmt.Printf(" %-20s settings from %s\n", module, layer.From)
+21 -26
View File
@@ -4,7 +4,6 @@ import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"flag"
"fmt"
@@ -283,10 +282,10 @@ func pushCommand(ctx context.Context, args []string) error {
asked = append(asked, n.Name)
}
sending, refusals := composeEach(asked, func(node string) ([]map[string]any, error) {
sending, refusals := composeEach(asked, func(node string) (sendable, error) {
plan, settings, err := planFor(ctx, open, node)
if err != nil {
return nil, err
return sendable{}, err
}
// A module assigned here that this machine cannot host is said and left out, not fatal: the
// healthy modules beside it are still resolved and sent. Reported so it is not silently
@@ -300,7 +299,7 @@ func pushCommand(ctx context.Context, args []string) error {
sentDigest := map[string]string{}
for _, s := range sending {
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
body, err := s.declared.Body()
if err != nil {
return err
}
@@ -318,7 +317,7 @@ func pushCommand(ctx context.Context, args []string) error {
return err
}
sentDigest[s.node] = digest
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.resources))
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
}
fmt.Printf("\n%d node(s) told\n", len(sending))
@@ -374,17 +373,17 @@ func pushCommand(ctx context.Context, args []string) error {
// (novox/hq ADR 0066). The earlier cut routed these through sendTo, which is
// all-or-nothing — so one swept machine's compose error failed the operator's named
// push and skipped its --wait, the very intolerance the main path exists to avoid.
sending, refused := composeEach(also, func(node string) ([]map[string]any, error) {
sending, refused := composeEach(also, func(node string) (sendable, error) {
plan, settings, err := planFor(ctx, open, node)
if err != nil {
return nil, err
return sendable{}, err
}
reportUnhostable(node, plan)
return declarationWith(ctx, open, node, plan, settings, gens, Allocating)
})
refusals = append(refusals, refused...)
for _, s := range sending {
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
body, err := s.declared.Body()
if err != nil {
return err
}
@@ -398,7 +397,7 @@ func pushCommand(ctx context.Context, args []string) error {
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
return err
}
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.resources))
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
}
// Every candidate this round is marked handled — the sent ones so they are not
// re-listed, and the refused ones so a machine that cannot be composed does not make
@@ -458,8 +457,8 @@ func waitForApplied(ctx context.Context, inv *inventory.Inventory, node, digest
// readyNode is one machine and the declaration it would be sent.
type readyNode struct {
node string
resources []map[string]any
node string
declared sendable
}
// composeEach works out what each named machine should be, and never lets one machine's answer
@@ -480,21 +479,21 @@ type readyNode struct {
// The all-or-nothing rule is kept where it means something — sendTo, which rotates a credential
// across two machines that must agree — and dropped here, where it never did.
func composeEach(names []string,
compose func(node string) ([]map[string]any, error)) ([]readyNode, []string) {
compose func(node string) (sendable, error)) ([]readyNode, []string) {
var sending []readyNode
var refusals []string
for _, name := range names {
resources, err := compose(name)
declared, err := compose(name)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue
}
if len(resources) == 0 {
if len(declared.Resources) == 0 {
fmt.Printf("%s is assigned nothing — skipped\n", name)
continue
}
sending = append(sending, readyNode{name, resources})
sending = append(sending, readyNode{name, declared})
}
return sending, refusals
}
@@ -533,11 +532,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
return err
}
type ready struct {
node string
resources []map[string]any
}
var sending []ready
var sending []readyNode
var refusals []string
for _, name := range names {
plan, settings, err := planFor(ctx, open, name)
@@ -546,12 +541,12 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
continue
}
reportUnhostable(name, plan)
resources, err := declarationWith(ctx, open, name, plan, settings, gens, Allocating)
declared, err := declarationWith(ctx, open, name, plan, settings, gens, Allocating)
if err != nil {
refusals = append(refusals, fmt.Sprintf("%s:\n%v", name, err))
continue
}
sending = append(sending, ready{name, resources})
sending = append(sending, readyNode{name, declared})
}
if len(refusals) > 0 {
return fmt.Errorf("nothing was sent. %d machine(s) could not be resolved:\n\n%s",
@@ -565,7 +560,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
defer server.Close()
for _, s := range sending {
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
body, err := s.declared.Body()
if err != nil {
return err
}
@@ -579,7 +574,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
return err
}
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.resources))
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
}
return nil
}
@@ -610,11 +605,11 @@ func wouldSend(ctx context.Context, open *stores,
if err != nil {
continue
}
resources, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
declared, err := declarationWith(ctx, open, n.Name, plan, settings, gens, Reading)
if err != nil {
continue
}
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": resources})
body, err := declared.Body()
if err != nil {
return nil, err
}
+4 -4
View File
@@ -18,11 +18,11 @@ import (
func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) {
sending, refusals := composeEach(
[]string{"anchor", "home-server", "laptop"},
func(node string) ([]map[string]any, error) {
func(node string) (sendable, error) {
if node == "anchor" {
return nil, errors.New(`nothing provides "acme-ca", wanted by route-proxy`)
return sendable{}, errors.New(`nothing provides "acme-ca", wanted by route-proxy`)
}
return []map[string]any{{"id": node + ".thing"}}, nil
return sendable{Resources: []map[string]any{{"id": node + ".thing"}}}, nil
})
var told []string
@@ -42,7 +42,7 @@ func TestOneUnresolvableNodeStillLetsTheRestBeSent(t *testing.T) {
// And a machine assigned nothing is neither sent nor a refusal — it is nothing to say.
func TestAMachineAssignedNothingIsNotARefusal(t *testing.T) {
sending, refusals := composeEach([]string{"spare"},
func(string) ([]map[string]any, error) { return nil, nil })
func(string) (sendable, error) { return sendable{}, nil })
if len(sending) != 0 || len(refusals) != 0 {
t.Errorf("a machine assigned nothing was treated as something: %v / %v", sending, refusals)
}
+96
View File
@@ -0,0 +1,96 @@
package main
import (
"context"
"encoding/json"
"sort"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// sendable is a declaration as a machine is sent it: its resources and, for an adopted node, its
// mode and which modules were taken on it (novox/hq ADR 0100).
//
// **Body is the only place the envelope is marshalled.** It was written by hand at every send site,
// in the digest the mesh compares, and in `plan --json`; a key added at one and not another would
// make a machine look out of date for ever, or send something `plan` never showed.
type sendable struct {
Resources []map[string]any
// Adoption is nil for a converged node, and then the body is byte for byte what it was before
// adoption existed: an older host parses the envelope strictly and would refuse the key.
Adoption *adoptionEnvelope
}
// adoptionEnvelope is what an adopted node is told about its mode. Taken is every module taken on
// it that it runs; Untaken is, for every module it runs that is not taken, the ids of that module's
// files and containers — the resources the host keeps as found until the module is taken. Ids
// rather than a rule to split them by, because a module's name may contain a dot.
type adoptionEnvelope struct {
Taken []string `json:"taken"`
Untaken map[string][]string `json:"untaken,omitempty"`
}
// Body is the declaration's bytes, as sent and as digested.
func (s sendable) Body() ([]byte, error) {
envelope := map[string]any{"declaration": 1, "resources": s.Resources}
if s.Adoption != nil {
envelope["adoption"] = s.Adoption
}
return json.Marshal(envelope)
}
// adoptionOf is the envelope for a node, nil when it is converged.
func adoptionOf(ctx context.Context, inv *inventory.Inventory, node string,
plan catalogue.Resolution, composed catalogue.Composed) (*adoptionEnvelope, error) {
record, err := inv.NodeByName(ctx, node)
if err != nil {
return nil, err
}
if !record.Adopted {
return nil, nil
}
taken, err := inv.Taken(ctx, node)
if err != nil {
return nil, err
}
return adoptionFor(plan, taken, composed), nil
}
// adoptionFor is the envelope computed from what was taken and who owns each resource.
//
// Every module the node runs that is not taken is untaken — including one pulled in by another
// rather than assigned: what is found is kept until its module is taken, whoever put it there.
func adoptionFor(plan catalogue.Resolution, taken []string,
composed catalogue.Composed) *adoptionEnvelope {
isTaken := map[string]bool{}
for _, m := range taken {
isTaken[m] = true
}
out := &adoptionEnvelope{Taken: []string{}}
runs := map[string]bool{}
for _, m := range plan.Modules {
runs[m.Module] = true
if isTaken[m.Module] {
out.Taken = append(out.Taken, m.Module)
}
}
sort.Strings(out.Taken)
for _, r := range composed.Resources {
id, _ := r["id"].(string)
module, owned := composed.Owner[id]
if !owned || !runs[module] || isTaken[module] {
continue
}
switch r["type"] {
case "file", "container":
default:
continue
}
if out.Untaken == nil {
out.Untaken = map[string][]string{}
}
out.Untaken[module] = append(out.Untaken[module], id)
}
return out
}
+170
View File
@@ -0,0 +1,170 @@
package main
import (
"bytes"
"encoding/json"
"io"
"os"
"reflect"
"strings"
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// novox/hq ADR 0100: every declaration an adopted node is sent says it is adopted and which modules
// were taken on it; a converged node's declaration is byte for byte what it was.
// helloWeb is a module the predecessor also runs: a page and a server under names it uses.
func helloWeb() catalogue.Manifest {
return catalogue.Manifest{Module: "hello-web", Version: "1",
Resources: []map[string]any{
{"id": "page", "type": "file", "path": "/var/lib/hello-web/index.html", "content": "hello"},
{"id": "server", "type": "container", "name": "hello-web",
"image": "registry.example/hello@sha256:" + strings.Repeat("a", 64)},
{"id": "served", "type": "directory", "path": "/var/lib/hello-web"},
}}
}
// composed is what node would be sent now, as push composes it.
func composed(t *testing.T, open *stores, node string) sendable {
t.Helper()
plan, settings, err := planFor(t.Context(), open, node)
if err != nil {
t.Fatal(err)
}
declared, err := declarationFor(t.Context(), open, node, plan, settings)
if err != nil {
t.Fatal(err)
}
return declared
}
func TestAConvergedDeclarationIsByteForByteWhatItWas(t *testing.T) {
open := aMesh(t)
register(t, open, helloWeb())
if _, err := assign(t.Context(), open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
declared := composed(t, open, "anchor")
if declared.Adoption != nil {
t.Fatal("a converged node was given an adoption envelope")
}
body, err := declared.Body()
if err != nil {
t.Fatal(err)
}
// The envelope exactly as every send site marshalled it before adoption existed.
before, err := json.Marshal(map[string]any{"declaration": 1, "resources": declared.Resources})
if err != nil {
t.Fatal(err)
}
if !bytes.Equal(body, before) {
t.Fatalf("a converged declaration changed:\n%s\n%s", body, before)
}
if bytes.Contains(body, []byte(`"adoption"`)) {
t.Fatal("a converged declaration names adoption; an older host would refuse it")
}
}
func TestAnAdoptedDeclarationCarriesItsModeAndWhatWasTaken(t *testing.T) {
open := aMesh(t)
ctx := t.Context()
register(t, open, helloWeb())
if err := open.inventory.SetAdopted(ctx, "anchor", true); err != nil {
t.Fatal(err)
}
if _, err := assign(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
declared := composed(t, open, "anchor")
if declared.Adoption == nil {
t.Fatal("an adopted node's declaration does not say it is adopted")
}
untaken := declared.Adoption.Untaken["hello-web"]
if !reflect.DeepEqual(untaken, []string{"hello-web.page", "hello-web.server"}) {
t.Fatalf("hello-web's files and containers are not named untaken: %v", declared.Adoption)
}
if len(declared.Adoption.Taken) != 0 {
t.Fatalf("nothing was taken, and the declaration says %v", declared.Adoption.Taken)
}
body, err := declared.Body()
if err != nil {
t.Fatal(err)
}
var envelope map[string]any
if err := json.Unmarshal(body, &envelope); err != nil {
t.Fatal(err)
}
adoption, _ := envelope["adoption"].(map[string]any)
if _, ok := adoption["taken"].([]any); !ok {
t.Fatalf("taken is not a list on the wire, even empty: %s", body)
}
// The digest the mesh compares is the digest of what is sent: status and push agree.
would, err := wouldSend(ctx, open, mustNodes(t, open))
if err != nil {
t.Fatal(err)
}
if would["anchor"] != digestOf(body) {
t.Fatal("the digest the mesh compares is not of the declaration push sends")
}
// plan --json prints that same envelope.
printed := stdoutOf(t, func() error { return planCommand(ctx, []string{"anchor", "--json"}) })
var compact bytes.Buffer
if err := json.Compact(&compact, []byte(printed)); err != nil {
t.Fatalf("plan --json is not JSON: %v\n%s", err, printed)
}
if digestOf(compact.Bytes()) != digestOf(body) {
t.Fatalf("plan --json shows something other than what push sends:\n%s", printed)
}
// Taking the module moves its resources out of untaken.
if err := open.inventory.Take(ctx, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
declared = composed(t, open, "anchor")
if _, still := declared.Adoption.Untaken["hello-web"]; still {
t.Fatalf("a taken module is still untaken: %v", declared.Adoption)
}
if !reflect.DeepEqual(declared.Adoption.Taken, []string{"hello-web"}) {
t.Fatalf("taken is %v", declared.Adoption.Taken)
}
}
func mustNodes(t *testing.T, open *stores) []inventory.Node {
t.Helper()
nodes, err := open.inventory.Nodes(t.Context())
if err != nil {
t.Fatal(err)
}
return nodes
}
// stdoutOf is what run printed.
func stdoutOf(t *testing.T, run func() error) string {
t.Helper()
r, w, err := os.Pipe()
if err != nil {
t.Fatal(err)
}
saved := os.Stdout
os.Stdout = w
done := make(chan string)
go func() {
all, _ := io.ReadAll(r)
done <- string(all)
}()
runErr := run()
os.Stdout = saved
w.Close()
out := <-done
if runErr != nil {
t.Fatalf("%v\n%s", runErr, out)
}
return out
}
+30
View File
@@ -146,6 +146,34 @@ func (r Rendering) machinePort(module string, wanted int) int {
// both call something "config", and without this the second would silently replace the first —
// the node applying one of them and reporting success.
func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) {
composed, err := r.Compose(with)
if err != nil {
return nil, err
}
return composed.Resources, nil
}
// Composed is a declaration's resources and which module each came from.
//
// Owner is kept beside the resources because a resource id cannot be split back into its module:
// a module's name may itself contain a dot. What the mesh adds of its own — an opening, the guard —
// has no owner.
type Composed struct {
Resources []map[string]any
Owner map[string]string
}
// Compose is Declaration with the owner of every resource said.
func (r Resolution) Compose(with Rendering) (Composed, error) {
owner := map[string]string{}
resources, err := r.compose(with, owner)
if err != nil {
return Composed{}, err
}
return Composed{Resources: resources, Owner: owner}, nil
}
func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[string]any, error) {
// Where each provision's credentials land, so a contribution can name the file rather than
// carry a value the mesh does not have.
directories := map[string]string{}
@@ -521,6 +549,7 @@ func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) {
if renamed := reflectsRenamed(m.Module, resource["restart-on"]); renamed != nil {
copied["restart-on"] = renamed
}
owner[fmt.Sprint(copied["id"])] = m.Module
out = append(out, copied)
}
@@ -534,6 +563,7 @@ func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) {
}
for _, fact := range given {
fact["id"] = m.Module + "." + fmt.Sprint(fact["id"])
owner[fmt.Sprint(fact["id"])] = m.Module
out = append(out, fact)
}
}