Compare commits

..
Author SHA1 Message Date
jschoubben bf204f90f3 The console's build tool has the command's three shapes: a repository, a base (--on), everything behind
Rebuilding the forty-three modules that stand on the runtime image took forty-three tool calls
because the seat verb only knew a repository. `on` rebuilds every module built on a base, `behind`
rebuilds what is older than its source, both asked and not waited for, the daemon taking each result
in (issue 176). Nothing is required any more; a build naming nothing is refused by the command's usage.
2026-10-01 14:05:57 +02:00
62 changed files with 222 additions and 2899 deletions
+9 -9
View File
@@ -109,10 +109,10 @@ var (
func TestTakingAModuleNotOnTheNodeIsRefused(t *testing.T) {
open, _ := anAdoptedAnchor(t)
if _, err := take(t.Context(), open, "anchor", "nftables", takeOptions{Yes: true}); !errors.Is(err, inventory.ErrNotAssigned) {
if _, err := take(t.Context(), open, "anchor", "nftables"); !errors.Is(err, inventory.ErrNotAssigned) {
t.Fatalf("taking an unassigned module gave %v", err)
}
if _, err := take(t.Context(), open, "laptop", "network", takeOptions{Yes: true}); !errors.Is(err, inventory.ErrNotAdopted) {
if _, err := take(t.Context(), open, "laptop", "network"); !errors.Is(err, inventory.ErrNotAdopted) {
t.Fatalf("taking on a converged node gave %v", err)
}
}
@@ -143,7 +143,7 @@ func TestTheFlipIsRefusedWhileAFoundContainerIsHeld(t *testing.T) {
func TestTakingNamesWhatItReplaces(t *testing.T) {
open, _ := anAdoptedAnchor(t)
reportsHolding(t, open, heldContainer, heldFile)
said, err := take(t.Context(), open, "anchor", "hello-web", takeOptions{Yes: true})
said, err := take(t.Context(), open, "anchor", "hello-web")
if err != nil {
t.Fatal(err)
}
@@ -156,7 +156,7 @@ func TestTakingNamesWhatItReplaces(t *testing.T) {
func TestConvergingPreviewsThenChangesAndAdoptingKeepsWhatWasTaken(t *testing.T) {
open, sent := anAdoptedAnchor(t)
ctx := t.Context()
if _, err := take(ctx, open, "anchor", "hello-web", takeOptions{Yes: true}); err != nil {
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
reportsHolding(t, open, heldFile)
@@ -330,7 +330,7 @@ func digestIn(t *testing.T, preview string) string {
func TestTheFlipActsOnlyOnThePreviewTheOperatorSaw(t *testing.T) {
open, sent := anAdoptedAnchor(t)
ctx := t.Context()
if _, err := take(ctx, open, "anchor", "hello-web", takeOptions{Yes: true}); err != nil {
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
reportsHolding(t, open, heldFile)
@@ -396,7 +396,7 @@ func TestTheFlipActsOnlyOnThePreviewTheOperatorSaw(t *testing.T) {
func TestTheFlipHoldsTheNodeWhileItSends(t *testing.T) {
open, _ := anAdoptedAnchor(t)
ctx := t.Context()
if _, err := take(ctx, open, "anchor", "hello-web", takeOptions{Yes: true}); err != nil {
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
reportsHolding(t, open, heldFile)
@@ -441,7 +441,7 @@ func TestTheFlipHoldsTheNodeWhileItSends(t *testing.T) {
func TestThePreviewNamesEveryHeldKind(t *testing.T) {
open, _ := anAdoptedAnchor(t)
ctx := t.Context()
if _, err := take(ctx, open, "anchor", "hello-web", takeOptions{Yes: true}); err != nil {
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
since := time.Now()
@@ -484,7 +484,7 @@ func TestThePreviewNamesEveryHeldKind(t *testing.T) {
func TestTheFlipIsRefusedOnAnAccountNamingNothingReachable(t *testing.T) {
open, sent := anAdoptedAnchor(t)
ctx := t.Context()
if _, err := take(ctx, open, "anchor", "hello-web", takeOptions{Yes: true}); err != nil {
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
// Only a loopback listener: nothing off the machine, which is the same silence.
@@ -512,7 +512,7 @@ func TestTheFlipIsRefusedOnAnAccountNamingNothingReachable(t *testing.T) {
func TestAssigningWaitsForWhateverIsConvergingTheNode(t *testing.T) {
open, _ := anAdoptedAnchor(t)
ctx := t.Context()
if _, err := take(ctx, open, "anchor", "hello-web", takeOptions{Yes: true}); err != nil {
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
t.Fatal(err)
}
reportsHolding(t, open, heldFile)
+17 -206
View File
@@ -24,11 +24,6 @@ import (
func showMode(ctx context.Context, inv *inventory.Inventory, node inventory.Node) error {
if !node.Adopted {
fmt.Printf(" mode converged\n")
// A converged machine holds nothing, and can still run what nobody asked for
// (novox/hq ADR 0163): what it reports as strays is said whatever its mode.
if said, err := inv.AdoptionOf(ctx, node.Name); err == nil && len(said.Strays) > 0 {
showStrays(said.Strays)
}
return nil
}
fmt.Printf(" mode adopted since %s\n",
@@ -67,26 +62,11 @@ func showMode(ctx context.Context, inv *inventory.Inventory, node inventory.Node
if h.Kept != "" {
fmt.Printf(" %-17s original kept at %s\n", "", h.Kept)
}
for _, f := range comparisonLines(h) {
fmt.Printf(" %-17s %s\n", "", f)
}
}
showStrays(said.Strays)
fmt.Printf(" as of %s\n", said.At.Local().Format(time.DateTime))
return nil
}
// showStrays says what a machine runs that the mesh neither wrote nor holds (ADR 0163).
func showStrays(strays []inventory.Stray) {
if len(strays) == 0 {
return
}
fmt.Printf(" strays %d container(s) the mesh neither wrote nor holds:\n", len(strays))
for _, s := range strays {
fmt.Printf(" %-17s %s (%s)\n", "", s.Name, s.Detail)
}
}
// showTunnel is the node show lines about the tunnel an adopted node found and carried (novox/hq
// ADR 0105): what it presented at enrolment, and what it last said about taking it over.
func showTunnel(ctx context.Context, inv *inventory.Inventory, name string) error {
@@ -156,14 +136,7 @@ const DefaultFilter = "nftables"
// take is a module's cutover on an adopted node: the operator's act, done when that module's data
// has moved. From the next push its resources converge there like any other, replacing what the
// node found and holds for it.
// takeOptions is what a take was told about the differences it may pass (novox/hq ADR 0163).
type takeOptions struct {
Yes bool
Downgrade bool
Replace map[string]bool
}
func take(ctx context.Context, open *stores, node, module string, opts takeOptions) (string, error) {
func take(ctx context.Context, open *stores, node, module string) (string, error) {
inv := open.inventory
assigned, err := inv.Assigned(ctx, node)
if err != nil {
@@ -177,170 +150,27 @@ func take(ctx context.Context, open *stores, node, module string, opts takeOptio
}
}
}
// The comparison first (novox/hq ADR 0163): every held thing the module would replace, beside
// what the module declares, and the differences that refuse unless named.
reported, err := inv.AdoptionOf(ctx, node)
if err != nil {
return "", err
}
preview, refusals := comparisonOf(reported.Held, module, opts)
if len(refusals) > 0 {
return "", fmt.Errorf("taking %s on %s is refused:\n %s\n%s", module, node,
strings.Join(refusals, "\n "), preview)
}
if !opts.Yes {
return preview + fmt.Sprintf("\nnothing taken; `take %s %s --yes` cuts it over as previewed", node, module), nil
}
if err := inv.Take(ctx, node, module); err != nil {
return "", err
}
said := fmt.Sprintf("%s is taken on %s", module, node)
if preview != "" {
said += "; the next push replaces what the node found and holds for it:\n" + preview
reported, err := inv.AdoptionOf(ctx, node)
if err != nil {
return "", err
}
var replaces []string
for _, h := range reported.Held {
if h.Module == module {
replaces = append(replaces, " "+heldLine(h))
}
}
if len(replaces) > 0 {
said += "; the next push replaces what the node found and holds for it:\n" +
strings.Join(replaces, "\n")
}
return said + fmt.Sprintf("\n run `push %s` to cut it over", node), nil
}
// comparisonOf is a take's preview: for every held thing of the module, what runs beside what the
// module declares, and the refusals the differences earn unless the take named them
// (novox/hq ADR 0163): an image older than the one running, a declared file that differs from the
// found one. A narrowed port and a shared network are said and not refused.
func comparisonOf(held []inventory.Held, module string, opts takeOptions) (string, []string) {
var b strings.Builder
var refusals []string
for _, h := range held {
if h.Module != module {
continue
}
fmt.Fprintf(&b, " %s", heldLine(h))
if h.Kept != "" {
fmt.Fprintf(&b, ", original kept at %s", h.Kept)
}
b.WriteString("\n")
for _, line := range comparisonLines(h) {
fmt.Fprintf(&b, " %s\n", line)
}
f := factsOf(h)
if f.downgrade && !opts.Downgrade {
refusals = append(refusals, fmt.Sprintf("%s: the module's image (%s, made %s) is older than the one running (%s, made %s) — "+
"a service that migrated its data forward may not start on it; `--downgrade` to take it anyway",
h.Target, f.declaredImage, day(f.declaredCreated), f.image, day(f.imageCreated)))
}
if f.differs && !opts.Replace[h.Target] && !opts.Replace["*"] {
refusals = append(refusals, fmt.Sprintf("%s: the module's content differs from the file found; the lines above "+
"marked - are lost by taking it; `--replace %s` to replace it anyway, or declare the file partially",
h.Target, h.Target))
}
}
return b.String(), refusals
}
// facts is a held thing's facts as the preview reads them.
type facts struct {
image, imageCreated, declaredImage, declaredCreated string
downgrade, differs bool
networks map[string][]string
mounts, ports, declaredPorts, declaredVolumes []string
difference []string
}
func factsOf(h inventory.Held) facts {
var f facts
if h.Facts == nil {
return f
}
str := func(k string) string { s, _ := h.Facts[k].(string); return s }
list := func(k string) []string {
var out []string
if raw, ok := h.Facts[k].([]any); ok {
for _, x := range raw {
if s, ok := x.(string); ok {
out = append(out, s)
}
}
}
return out
}
f.image, f.imageCreated = str("image"), str("image_created")
f.declaredImage, f.declaredCreated = str("declared_image"), str("declared_image_created")
f.downgrade, _ = h.Facts["downgrade"].(bool)
f.differs, _ = h.Facts["differs"].(bool)
f.mounts, f.ports = list("mounts"), list("ports")
f.declaredPorts, f.declaredVolumes, f.difference = list("declared_ports"), list("declared_volumes"), list("difference")
if raw, ok := h.Facts["networks"].(map[string]any); ok {
f.networks = map[string][]string{}
for name, members := range raw {
var out []string
if ms, ok := members.([]any); ok {
for _, m := range ms {
if s, ok := m.(string); ok {
out = append(out, s)
}
}
}
f.networks[name] = out
}
}
return f
}
// comparisonLines says a held thing's facts the way a person weighs them.
func comparisonLines(h inventory.Held) []string {
f := factsOf(h)
var out []string
if f.image != "" || f.declaredImage != "" {
line := fmt.Sprintf("runs %s", orNone(f.image))
if f.imageCreated != "" {
line += " (made " + day(f.imageCreated) + ")"
}
line += "; the module declares " + orNone(f.declaredImage)
switch {
case f.declaredCreated != "":
line += " (made " + day(f.declaredCreated) + ")"
case f.declaredImage != "":
line += " (not on the machine yet, so its age is unknown)"
}
if f.downgrade {
line += " — DOWNGRADE"
}
out = append(out, line)
}
names := make([]string, 0, len(f.networks))
for n := range f.networks {
names = append(names, n)
}
sort.Strings(names)
for _, n := range names {
if members := f.networks[n]; len(members) > 0 {
out = append(out, fmt.Sprintf("on the network %s with %s, which may reach it by name and will not once it moves to the module's own network",
n, strings.Join(members, ", ")))
}
}
if len(f.ports) > 0 || len(f.declaredPorts) > 0 {
out = append(out, fmt.Sprintf("publishes %s; the module declares %s",
orNone(strings.Join(f.ports, " ")), orNone(strings.Join(f.declaredPorts, " "))))
}
if len(f.mounts) > 0 || len(f.declaredVolumes) > 0 {
out = append(out, fmt.Sprintf("mounts %s; the module declares %s",
orNone(strings.Join(f.mounts, " ")), orNone(strings.Join(f.declaredVolumes, " "))))
}
if f.differs {
out = append(out, "the declared content differs from the file found (- lost, + new):")
for _, d := range f.difference {
out = append(out, " "+d)
}
}
return out
}
// day is a timestamp as a person reads it in a preview: its date.
func day(stamp string) string {
if len(stamp) >= 10 {
return stamp[:10]
}
return stamp
}
// reportFreshFor is how old a node's account of itself may be for the flip to act on it. A
// variable so a test can age a report without waiting.
var reportFreshFor = 15 * time.Minute
@@ -766,31 +596,12 @@ func adopt(ctx context.Context, open *stores, node string) (string, error) {
// takeCommand, convergeCommand and adoptCommand are the command line's adapters to the acts above.
func takeCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("take", flag.ContinueOnError)
yes := set.Bool("yes", false, "cut over as previewed; without it the comparison is printed and nothing is taken")
downgrade := set.Bool("downgrade", false, "take it although the module's image is older than the one running")
var replace stringList
set.Var(&replace, "replace", "a found file's path whose content the module may replace although it differs (repeatable; * for every one)")
positionals, err := parseAround(set, args)
if err != nil {
return err
if len(args) != 2 {
return errors.New("take <node> <module>")
}
if len(positionals) != 2 {
return errors.New("take <node> <module> [--yes] [--downgrade] [--replace <path>]...")
}
opts := takeOptions{Yes: *yes, Downgrade: *downgrade, Replace: map[string]bool{}}
for _, r := range replace {
opts.Replace[r] = true
}
return runAct(ctx, func(open *stores) (string, error) { return take(ctx, open, positionals[0], positionals[1], opts) })
return runAct(ctx, func(open *stores) (string, error) { return take(ctx, open, args[0], args[1]) })
}
// stringList is a repeatable flag.
type stringList []string
func (l *stringList) String() string { return strings.Join(*l, ",") }
func (l *stringList) Set(v string) error { *l = append(*l, v); return nil }
func convergeCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("converge", flag.ContinueOnError)
yes := set.String("yes", "", "do it, naming the digest the preview printed; without it, only "+
+1 -1
View File
@@ -102,7 +102,7 @@ func commands(who Authenticator) http.Handler {
}))
// Adoption (novox/hq ADR 0100): the same acts as `take`, `converge` and `adopt`.
mux.HandleFunc("POST /take", acting(who, true, func(ctx context.Context, open *stores, in request) (string, error) {
return take(ctx, open, in.Node, in.Module, takeOptions{Yes: true})
return take(ctx, open, in.Node, in.Module)
}))
mux.HandleFunc("POST /converge", acting(who, false, func(ctx context.Context, open *stores, in request) (string, error) {
return converge(ctx, open, in.Node, in.Yes, in.Digest, in.Filter)
-17
View File
@@ -469,25 +469,10 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
}
fmt.Printf("\n%s %s, built on %s from %s\n",
manifest.Module, manifest.Version, result.On, short(result.Commit))
saysWhenThePolicyActs(ctx, open.inventory, manifest.Module)
fmt.Printf(" run `assign <node> %s` to put it somewhere\n", manifest.Module)
return nil
}
// saysWhenThePolicyActs tells whoever built a module that its upgrade policy will send the
// result on at once (novox/hq issue 126, ADR 0163): a person choreographing a data move must
// know which module will not wait for them.
func saysWhenThePolicyActs(ctx context.Context, inv *inventory.Inventory, module string) {
if u, err := inv.UpgradeOf(ctx, module); err == nil && u.RollOut {
how := "one machine at a time"
if u.Together {
how = "every machine at once"
}
fmt.Printf(" %s rolls out on build: the machines running it are sent this now, %s — "+
"`upgrade %s record` first if something must move before it does\n", module, how, module)
}
}
// takeIn is what the mesh does with a build's outcome, whoever hears it: the waiting command and
// the daemon that follows the role's events both come here (novox/hq issue 176), so a build's
// result reaches the catalogue whether or not the asker was still listening.
@@ -596,8 +581,6 @@ type answers struct {
// pair that answers "has it caught up", which waiting alone cannot (the sent digest is
// recorded at send, not at apply).
reported []inventory.Reported
// plans is what the last merges produced and where each stands (novox/hq ADR 0162).
plans []inventory.Plan
// refused is why a machine cannot be worked out at all, by name. A different thing from every
// other answer here: those are about a machine that was told something, and this is about one
// that cannot be told anything — it never reaches waiting, because nothing was computed for it
-52
View File
@@ -1,52 +0,0 @@
package main
import (
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/inventory"
)
// A merge that rebuilds a base rebuilds what stands on it, through every layer, and nothing else
// (novox/hq issue 186): the runtime image moving means every module built on it moves too, and a
// module built on one of those moves as well.
func TestAMergeOfABaseTakesWhatStandsOnItAlong(t *testing.T) {
entry := func(name string) inventory.Entry {
return inventory.Entry{Manifest: catalogue.Manifest{Module: name}}
}
entries := []inventory.Entry{entry("mesh-tools"), entry("shop"), entry("shop-plugin"), entry("postgres"), entry("unrelated")}
against := map[string][]string{
"shop": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
"shop-plugin": {catalogue.ArtifactStoreScheme + "shop/runtime@sha256:b"},
"postgres": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
"unrelated": {catalogue.ArtifactStoreScheme + "alpine/base@sha256:c"},
}
got := dependentsOf([]inventory.Entry{entry("mesh-tools")}, entries, against)
var names []string
for _, e := range got {
names = append(names, e.Manifest.Module)
}
want := map[string]bool{"shop": true, "shop-plugin": true, "postgres": true}
if len(names) != len(want) {
t.Fatalf("rebuilt %v; wanted exactly the three that stand on the runtime, directly or through shop", names)
}
for _, n := range names {
if !want[n] {
t.Fatalf("%s was rebuilt and stands on nothing that moved (%v)", n, names)
}
}
// The dependents come in base order when the merge orders them: the runtime, then shop, then
// the plugin that stands on shop.
ordered := orderByBases(append([]inventory.Entry{entry("mesh-tools")}, got...), against)
pos := map[string]int{}
for i, e := range ordered {
pos[e.Manifest.Module] = i
}
if !(pos["mesh-tools"] < pos["shop"] && pos["shop"] < pos["shop-plugin"]) {
t.Fatalf("not in base order: %v", ordered)
}
// Nothing moved: nothing follows.
if more := dependentsOf(nil, entries, against); len(more) != 0 {
t.Fatalf("with nothing moved, %d module(s) were rebuilt", len(more))
}
}
+1 -14
View File
@@ -72,8 +72,6 @@ func run() error {
return askCommand(ctx, args[1:])
case "builds":
return buildsCommand(ctx, args[1:])
case "plans":
return plansCommand(ctx, args[1:])
case "pin":
return pinCommand(ctx, args[1:], true)
case "unpin":
@@ -205,8 +203,7 @@ func usage() {
licence refresh <name> mint a new access token and seal it to every holder
rotate <provision> [--consumer <n>] a new credential for every holder, both ends at once
ask <module> <tool> [json] call one of a module's tools over the broker, and print its answer
pin <node> <provision> <from-node> <module>
which provider this one gets a provision from: the module, and its node
pin <node> <provision> <from> which node this one gets a provision from
unpin <node> <provision> put that question back
plan <node> [--files|--json] what that node would run, and why
push [<node>] [--behind] send a node everything it should be, or only those that need it
@@ -255,22 +252,12 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error {
switch {
case err != nil && result.Failed != "":
fmt.Printf("%s: %v\n", result.ID, err)
if result.Module != "" {
planBuilt(ctx, b.open, result.Module, result.Commit, result.Failed)
} else {
planFailedBuild(ctx, b.open, result)
}
return nil
case err != nil:
fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err)
if manifest.Module != "" {
planBuilt(ctx, b.open, manifest.Module, result.Commit, err.Error())
}
return nil
}
fmt.Printf("%s: %s %s registered, built on %s from %s\n",
result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit))
saysWhenThePolicyActs(ctx, b.inv, manifest.Module)
planBuilt(ctx, b.open, manifest.Module, result.Commit, "")
return nil
}
+12 -9
View File
@@ -411,8 +411,8 @@ func describeOffers(offers []catalogue.Offer) string {
// database should not change where an existing machine gets its data the day a second one
// arrives.
func pinCommand(ctx context.Context, args []string, setting bool) error {
if setting && len(args) != 4 {
return errors.New("pin <node> <provision> <from-node> <module>")
if setting && len(args) != 3 {
return errors.New("pin <node> <provision> <from-node>")
}
if !setting && len(args) != 2 {
return errors.New("unpin <node> <provision>")
@@ -431,12 +431,17 @@ func pinCommand(ctx context.Context, args []string, setting bool) error {
fmt.Printf("%s is no longer told where to get %s from\n", args[0], args[1])
return nil
}
// The provider's node may be this same machine: two modules beside the consumer can both
// answer a provision, and then the module is the whole question (novox/hq #258).
if err := inv.PinProvision(ctx, args[0], args[1], args[2], args[3]); err != nil {
if args[0] == args[2] {
// Allowed by nothing here, and worth saying rather than resolving into a confusing
// refusal later: a node providing something to itself is a node-scoped provision, and
// this field is for the other kind.
return fmt.Errorf("%s cannot get %s from itself; that would be a provision this machine "+
"provides, which does not need saying", args[0], args[1])
}
if err := inv.PinProvision(ctx, args[0], args[1], args[2]); err != nil {
return err
}
fmt.Printf("%s gets %s from %s/%s\n", args[0], args[1], args[2], args[3])
fmt.Printf("%s gets %s from %s\n", args[0], args[1], args[2])
fmt.Printf(" run `push %s` to send it\n", args[0])
return nil
}
@@ -711,9 +716,7 @@ func claimsFor(ctx context.Context, inv *inventory.Inventory, m catalogue.Manife
claimed := seatClaimed{Seat: c.Name, Scope: c.At()}
if s, known := byName[c.Name]; known {
claimed.Scope = s.Scope
// The verbs the runtime serves for the seat: the claim's own when it names them
// (ADR 0160), else every verb the seat promises, which its tools then answer.
claimed.Serves = c.ServesFor(catalogue.Manifest{Tools: catalogue.VerbNames(s.Serves)})
claimed.Serves = catalogue.VerbNames(s.Serves)
}
out = append(out, claimed)
}
-3
View File
@@ -567,9 +567,6 @@ func onTheNetwork(ctx context.Context, inv *inventory.Inventory,
catalogue.Node{Name: p.Name, Site: p.Site, Capabilities: caps},
catalogue.World{Unchecked: true, Holdings: holdings})
if err != nil {
// Said, not skipped in silence: a machine dropped here loses its address, and every
// plan that names it fails in another module's words (novox/hq issue 188).
fmt.Fprintf(os.Stderr, "%s is not counted as on the network: it does not resolve: %v\n", p.Name, err)
continue
}
for _, m := range got.Modules {
+3 -19
View File
@@ -7,7 +7,6 @@ import (
"errors"
"flag"
"fmt"
"os"
"sort"
"strings"
@@ -282,9 +281,8 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
for _, o := range others {
got, err := catalogue.Resolve(shelf, o.assigned, o.node, catalogue.World{Unchecked: true, Holdings: holdings})
if err != nil {
// Said, not skipped: a machine dropped here offers nothing and holds nothing as far
// as every other machine's plan can tell (novox/hq issue 188).
fmt.Fprintf(os.Stderr, "%s is left out of the rest of the mesh: it does not resolve: %v\n", o.node.Name, err)
// Their set does not resolve for some other reason. Not this node's problem to
// report, and nothing of theirs is running, so it offers nothing.
continue
}
firstHeld = append(firstHeld, got.Claims...)
@@ -315,22 +313,8 @@ func theRestOfTheMesh(ctx context.Context, inv *inventory.Inventory,
world := catalogue.World{Offered: offered, Held: firstHeld, Holdings: holdings}
var held []catalogue.Held
for _, o := range others {
// Each machine is resolved with its own pins, as its plan is: a machine that needs one to
// settle two providers would otherwise be refused here and vanish from the mesh — every
// seat it holds unheld, every build that needs one refused (2026-10-01, the control node;
// novox/hq issue 188).
theirs := world
if pins, err := inv.PinsFor(ctx, o.node.Name); err == nil {
theirs.Pinned = pins
}
got, err := catalogue.Resolve(shelf, o.assigned, o.node, theirs)
got, err := catalogue.Resolve(shelf, o.assigned, o.node, world)
if err != nil {
// Said only for the whole-mesh view. With one machine excluded, the others are
// resolved without its offers, and one that consumes them cannot resolve here by
// design — that is not the machine being dropped, it is the view being partial.
if exclude == "" {
fmt.Fprintf(os.Stderr, "%s is left out of the rest of the mesh: it does not resolve: %v\n", o.node.Name, err)
}
continue
}
held = append(held, got.Claims...)
+1 -59
View File
@@ -4,7 +4,6 @@ import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"flag"
"fmt"
@@ -113,10 +112,7 @@ func serve(ctx context.Context) error {
// while everything else about it looks correct.
// And build results nobody was waiting for. A build triggered any other way than `build`
// would otherwise be reported into the void, which is the same as not reporting it.
server.Records(builds{inv, open})
// Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for
// machines to report moves when they have, and a plan left by a replaced controller resumes.
go planTicker(ctx, open)
server.Records(builds{inv})
// And what the catalogue decided a build meant. The builder's own result is already handled
// above; this is the other half — the control plane is the only one of the three that knows
// which machines run the thing, so it is the one that acts (novox/hq ADR 0072).
@@ -387,15 +383,6 @@ func pushCommand(ctx context.Context, args []string) error {
}
release()
fmt.Printf("\n%d node(s) told\n", len(sending))
// And each machine's memberships, as every other send does (ADR 0160): a push is the one most
// operators run, and on 2026-10-01 it was the one path that issued none.
var told []string
for _, s := range sending {
told = append(told, s.node)
}
if err := issueMemberships(ctx, open, server, told); err != nil {
return err
}
// **A named push leaves the mesh consistent, not just the machine it named** (novox/hq
// issue 057, ADR 0083). Assigning a cross-node consumer mints a provision, and the PROVIDER's
@@ -703,51 +690,6 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
}
fmt.Printf(" sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
}
// And every assignment on those machines its membership (novox/hq ADR 0160): composed from the
// same records the bus's accounts are, so what a runtime serves and what its account may are one
// composition. Issued after the declaration, because the runtime it is for arrives with it.
return issueMemberships(ctx, open, server, names)
}
// issueMemberships publishes the membership of every module on the named machines.
func issueMemberships(ctx context.Context, open *stores, server *link.Server, names []string) error {
records, err := open.inventory.BusRecords(ctx)
if err != nil {
return err
}
where := broker.PlacementsOf(records, records.Interchangeable)
bus, ok := server.Bus().(link.OverNATS)
if !ok {
return nil
}
// The declarations are sent and recorded by now; a membership that cannot be issued is said
// and does not unsay them. Every runtime without one serves the shape it derives (ADR 0160), so
// the push stands, the first failure is named once, and the next push tries again.
issued, failed := 0, 0
var first error
for _, node := range names {
for _, d := range records.Assigned[node] {
body, err := json.Marshal(broker.MembershipFor(node, d, where))
if err != nil {
return err
}
if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil {
if first == nil {
first = err
}
failed++
continue
}
issued++
}
}
if issued > 0 {
fmt.Printf(" issued %d membership(s)\n", issued)
}
if failed > 0 {
fmt.Printf(" %d membership(s) could not be issued; the first: %v — the machines keep what "+
"they derive until the next push\n", failed, first)
}
return nil
}
+1 -4
View File
@@ -39,9 +39,6 @@ type meshStatus struct {
// whose is older is still working — and Waiting cannot tell those apart, because the sent
// digest is recorded at send, not at apply.
Reported []machineReported `json:"reported"`
// Plans is what the last merges produced and where each stands (novox/hq ADR 0162): the
// open ones first, each saying its tier, what it waits for, and whether it has waited too long.
Plans []planStatus `json:"plans"`
// Unresolved is every machine that cannot be worked out at all, with what the mesh said when
// it tried. **A machine here is in none of the lists above**: nothing was computed for it, so
// there is nothing to compare it against and nothing it can be behind — which is why a
@@ -155,7 +152,7 @@ func statusAsJSON(asked answers) ([]byte, error) {
out := meshStatus{Machines: len(nodes), Wrong: []machineDoing{},
Quiet: []machineQuiet{}, Behind: []moduleBehind{}, Waiting: []machineWaiting{},
Reported: []machineReported{}, Unresolved: []machineUnresolved{},
Network: asked.network, Adopted: adoptedNodes(nodes), Plans: planStatuses(asked.plans, time.Now())}
Network: asked.network, Adopted: adoptedNodes(nodes)}
// In a stated order, so two readings of an unchanged mesh are the same document.
untakenNodes := make([]string, 0, len(asked.untaken))
for name := range asked.untaken {
-742
View File
@@ -1,742 +0,0 @@
package main
import (
"context"
"flag"
"fmt"
"sort"
"strings"
"time"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A merge produces a tiered plan the mesh keeps (novox/hq ADR 0162).
//
// The handler that hears the merge computes the plan from the catalogue's one dependency relation,
// writes it to the store, asks the first tier and returns — the receive loop is never held by a
// build. Every outcome taken in advances the plan it belongs to; a ticker advances what outcomes
// alone cannot (a tier waiting for machines to report); a controller replaced mid-plan finds the
// plan where it left it.
// planWaitBound is how long a plan may wait on one thing before `status` names it red.
const planWaitBound = 30 * time.Minute
// tiersOf sorts a set of modules into tiers along the ordering edges among them: tier 0 depends
// on nothing else in the set, tier 1 only on tier 0, and so on. An edge to a module outside the set says
// nothing about the order inside it. A cycle — which the catalogue should never produce — puts
// what remains in one last tier rather than losing it, and is said by the caller.
func tiersOf(set []string, edges []inventory.Edge) [][]string {
in := map[string]bool{}
for _, m := range set {
in[m] = true
}
deps := map[string]map[string]bool{}
for _, m := range set {
deps[m] = map[string]bool{}
}
for _, e := range edges {
// A code dependency — B packages A's source — rebuilds B with A, in the same tier: B's
// build needs nothing of A's first. The other kinds order: stands-on and declared after
// the base is built, built-by after the build machine is built and running — except for
// what the build machine itself stands on. The runtime image is built by the builder and
// the builder is built on the runtime image; the image comes first, built by the builder
// that is running, which is the only one there could be.
if !in[e.From] || !in[e.To] || e.From == e.To || e.Kind == inventory.EdgePackages {
continue
}
if e.Kind == inventory.EdgeBuiltBy && isBaseOf(e.From, e.To, edges, in) {
continue
}
deps[e.From][e.To] = true
}
placed := map[string]bool{}
var tiers [][]string
for len(placed) < len(set) {
var tier []string
for _, m := range set {
if placed[m] {
continue
}
free := true
for d := range deps[m] {
if !placed[d] {
free = false
break
}
}
if free {
tier = append(tier, m)
}
}
if len(tier) == 0 {
// A cycle: everything left, together, and the caller says so.
for _, m := range set {
if !placed[m] {
tier = append(tier, m)
}
}
}
sort.Strings(tier)
for _, m := range tier {
placed[m] = true
}
tiers = append(tiers, tier)
}
return tiers
}
// isBaseOf says whether `to` stands on `base`, directly or through other bases in the set, along
// the build edges alone.
func isBaseOf(base, to string, edges []inventory.Edge, in map[string]bool) bool {
seen := map[string]bool{}
var walk func(string) bool
walk = func(m string) bool {
if m == base {
return true
}
if seen[m] {
return false
}
seen[m] = true
for _, e := range edges {
if e.From == m && in[e.To] && (e.Kind == inventory.EdgeStandsOn || e.Kind == inventory.EdgeDeclared) && walk(e.To) {
return true
}
}
return false
}
return walk(to)
}
// reachableFrom is the moved modules plus everything that depends on them, through every layer:
// what a merge rebuilds. Along the code and build edges only: a module *built by* the build machine
// is not changed by a new build machine, so a built-by edge orders and gates a plan and never
// widens it — the first plan of 2026-10-01 took the whole catalogue along for a controller change.
func reachableFrom(moved []string, edges []inventory.Edge) []string {
in := map[string]bool{}
for _, m := range moved {
in[m] = true
}
for grew := true; grew; {
grew = false
for _, e := range edges {
if e.Kind == inventory.EdgeBuiltBy {
continue
}
if in[e.To] && !in[e.From] {
in[e.From] = true
grew = true
}
}
}
out := make([]string, 0, len(in))
for m := range in {
out = append(out, m)
}
sort.Strings(out)
return out
}
// hasCycle says whether the tiers' last tier holds modules that still depend on each other.
func hasCycle(tiers [][]string, edges []inventory.Edge) bool {
if len(tiers) == 0 {
return false
}
last := map[string]bool{}
for _, m := range tiers[len(tiers)-1] {
last[m] = true
}
for _, e := range edges {
if last[e.From] && last[e.To] {
return true
}
}
return false
}
// planFor is the plan a merge produces: the moved modules and everything reachable from them,
// tiered, with the merge it answers.
func planOfMerge(m link.SourceMoved, moved []string, edges []inventory.Edge) inventory.Plan {
set := reachableFrom(moved, edges)
tiers := tiersOf(set, edges)
modules := map[string]*inventory.PlanModule{}
for _, name := range set {
modules[name] = &inventory.PlanModule{}
}
return inventory.Plan{
ID: fmt.Sprintf("plan-%d", time.Now().UnixNano()),
Repository: m.Owner + "/" + m.Repo,
Commit: m.Commit,
Created: time.Now().UTC(),
State: inventory.PlanBuilding,
Tiers: tiers,
Modules: modules,
}
}
// gates is what the next tier needs running from this one: a module of the tier that a later
// tier is built by — the runtime dependency — and whose policy rolls it out, must be applied by
// the machines running it before the next tier is asked. A base an image stands on need only be
// built; a source another module packages need not even be that.
func gates(p inventory.Plan, edges []inventory.Edge, rollsOut func(string) bool) []string {
if p.Tier >= len(p.Tiers) {
return nil
}
inTier := map[string]bool{}
all := map[string]bool{}
for _, tier := range p.Tiers {
for _, m := range tier {
all[m] = true
}
}
for _, m := range p.Tiers[p.Tier] {
inTier[m] = true
}
later := map[string]bool{}
for _, tier := range p.Tiers[p.Tier+1:] {
for _, m := range tier {
later[m] = true
}
}
seen := map[string]bool{}
var out []string
for _, e := range edges {
if later[e.From] && inTier[e.To] && e.Kind == inventory.EdgeBuiltBy && !seen[e.To] && rollsOut(e.To) &&
!isBaseOf(e.From, e.To, edges, all) {
seen[e.To] = true
out = append(out, e.To)
}
}
sort.Strings(out)
return out
}
// applied says whether every machine running the module has reported since the module was built.
func applied(module string, builtAt time.Time, running []string, reports []inventory.Reported) (bool, []string) {
at := map[string]*time.Time{}
for _, r := range reports {
at[r.Node] = r.At
}
var waiting []string
for _, n := range running {
if t := at[n]; t == nil || t.Before(builtAt) {
waiting = append(waiting, n)
}
}
return len(waiting) == 0, waiting
}
// askTier asks the build machine for every module of the tier, and marks each asked. A module
// the catalogue no longer holds, or whose ask could not be made, is a failure of the plan: a tier
// half asked is a tier that will never complete.
func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) error {
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
byName := map[string]inventory.Entry{}
for _, e := range entries {
byName[e.Manifest.Module] = e
}
now := time.Now().UTC()
for _, name := range p.Tiers[p.Tier] {
state := p.Modules[name]
if state == nil {
state = &inventory.PlanModule{}
p.Modules[name] = state
}
e, known := byName[name]
if !known {
state.State = "failed"
state.Why = "no longer in the catalogue"
p.State = inventory.PlanFailed
p.Note = name + " is no longer in the catalogue"
continue
}
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
fmt.Printf(" tier %d: ", p.Tier)
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 0); err != nil {
state.State = "failed"
state.Why = err.Error()
p.State = inventory.PlanFailed
p.Note = fmt.Sprintf("%s could not be asked for: %v", name, err)
continue
}
state.State = "asked"
state.AskedAt = &now
}
return nil
}
// planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and
// advances what that completes. Called from the daemon's take-in of every outcome.
func planBuilt(ctx context.Context, open *stores, module, commit, failed string) {
inv := open.inventory
plans, err := inv.OpenPlans(ctx)
if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err)
return
}
now := time.Now().UTC()
for i := range plans {
p := &plans[i]
if p.Tier >= len(p.Tiers) {
continue
}
inTier := false
for _, m := range p.Tiers[p.Tier] {
if m == module {
inTier = true
}
}
if !inTier {
continue
}
state := p.Modules[module]
if state == nil {
state = &inventory.PlanModule{}
p.Modules[module] = state
}
if failed != "" {
state.State = "failed"
state.Why = failed
p.State = inventory.PlanFailed
p.Note = fmt.Sprintf("%s failed to build in tier %d", module, p.Tier)
} else {
state.State = "built"
state.BuiltAt = &now
state.Commit = commit
}
if err := inv.SavePlan(ctx, *p); err != nil {
fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err)
continue
}
if p.State == inventory.PlanFailed {
fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note)
}
}
advancePlans(ctx, open)
}
// advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built
// and whose gates are applied gives way to the next; the last tier done is the plan done. Called
// after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes.
func advancePlans(ctx context.Context, open *stores) {
inv := open.inventory
plans, err := inv.OpenPlans(ctx)
if err != nil {
fmt.Printf("plans: cannot read them: %v\n", err)
return
}
if len(plans) == 0 {
return
}
edges, err := inv.Dependencies(ctx)
if err != nil {
fmt.Printf("plans: cannot read the dependencies: %v\n", err)
return
}
rollsOut := func(module string) bool {
u, err := inv.UpgradeOf(ctx, module)
return err == nil && u.RollOut
}
for i := range plans {
p := &plans[i]
for p.Open() {
moved, err := advanceOnce(ctx, open, p, edges, rollsOut)
if err != nil {
fmt.Printf("%s: %v\n", p.ID, err)
break
}
if err := inv.SavePlan(ctx, *p); err != nil {
fmt.Printf("%s: cannot keep the plan: %v\n", p.ID, err)
break
}
if !moved {
break
}
}
}
}
// advanceOnce takes one step of one plan and says whether anything changed.
func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
edges []inventory.Edge, rollsOut func(string) bool) (bool, error) {
inv := open.inventory
if p.Tier >= len(p.Tiers) {
p.State = inventory.PlanDone
fmt.Printf("%s: done — %s at %s, %d tier(s)\n", p.ID, p.Repository, short(p.Commit), len(p.Tiers))
return true, nil
}
tier := p.Tiers[p.Tier]
// Not yet asked: ask.
unasked := 0
for _, m := range tier {
if s := p.Modules[m]; s == nil || s.State == "" {
unasked++
}
}
if unasked == len(tier) {
if err := askTier(ctx, inv, p); err != nil {
return false, err
}
return true, nil
}
// Asked: wait for every build.
var latest time.Time
for _, m := range tier {
s := p.Modules[m]
if s == nil || s.State != "built" {
return false, nil
}
if s.BuiltAt != nil && s.BuiltAt.After(latest) {
latest = *s.BuiltAt
}
}
// Built: send every module of the tier whose policy rolls out, once, to the machines running
// it — whether or not its source commit moved. A dependent rebuilt because its base moved, or
// a module that packages another repository's source, keeps its commit; the catalogue announces
// no move for it and its machines would keep the old image until somebody pushed (novox/hq
// issue 189). A module whose policy records is built and left, as its policy says.
for _, m := range tier {
state := p.Modules[m]
if state == nil || state.SentAt != nil || !rollsOut(m) {
continue
}
running, err := inv.Running(ctx, m)
if err != nil {
return false, err
}
now := time.Now().UTC()
state.SentAt = &now
if len(running) == 0 {
continue
}
if err := sendTo(ctx, open, running); err != nil {
return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(running, ", "), p.Tier, err)
}
fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(running, ", "))
return true, nil
}
// And wait for what the next tier needs running.
needed := gates(*p, edges, rollsOut)
if len(needed) > 0 {
reports, err := inv.LastReports(ctx)
if err != nil {
return false, err
}
var waiting []string
for _, m := range needed {
running, err := inv.Running(ctx, m)
if err != nil {
return false, err
}
state := p.Modules[m]
if state == nil {
state = &inventory.PlanModule{}
p.Modules[m] = state
}
// The plan sends what it waits for. A rebuild from the same source commit is not a
// move the catalogue announces — the build machine rebuilt for a controller change
// is one — so the roll-out that opens this gate is the plan's to make, once, and
// the reports that open it are the ones after the send.
since := latest
if state.BuiltAt != nil {
since = *state.BuiltAt
}
if state.SentAt != nil && state.SentAt.After(since) {
since = *state.SentAt
}
if ok, on := applied(m, since, running, reports); !ok {
waiting = append(waiting, fmt.Sprintf("%s on %s", m, strings.Join(on, ", ")))
}
}
if len(waiting) > 0 {
note := "tier " + fmt.Sprint(p.Tier) + " built; waiting for " + strings.Join(waiting, "; ") + " to be applied"
changed := p.State != inventory.PlanRolling || p.Note != note
p.State = inventory.PlanRolling
p.Note = note
return changed, nil
}
}
p.Tier++
p.State = inventory.PlanBuilding
p.Note = ""
if p.Tier < len(p.Tiers) {
fmt.Printf("%s: tier %d done; asking tier %d: %s\n", p.ID, p.Tier-1, p.Tier, strings.Join(p.Tiers[p.Tier], ", "))
}
return true, nil
}
// planTicker advances open plans on a timer, for the steps outcomes alone cannot take.
func planTicker(ctx context.Context, open *stores) {
advancePlans(ctx, open)
tick := time.NewTicker(30 * time.Second)
defer tick.Stop()
for {
select {
case <-ctx.Done():
return
case <-tick.C:
advancePlans(ctx, open)
}
}
}
// planLine is one plan as `status` says it.
func planLine(p inventory.Plan, now time.Time) string {
where := fmt.Sprintf("tier %d of %d", min(p.Tier+1, len(p.Tiers)), len(p.Tiers))
switch p.State {
case inventory.PlanDone:
return fmt.Sprintf("%s %s done, %d tier(s)", p.Repository, short(p.Commit), len(p.Tiers))
case inventory.PlanFailed:
return fmt.Sprintf("%s %s FAILED at %s: %s", p.Repository, short(p.Commit), where, p.Note)
}
since := now.Sub(p.Updated).Round(time.Minute)
late := ""
if since > planWaitBound {
late = " — LATE"
}
what := "building"
if p.State == inventory.PlanRolling {
what = p.Note
}
return fmt.Sprintf("%s %s %s, %s for %s%s", p.Repository, short(p.Commit), where, what, since, late)
}
// planFailedBuild marks the module a failed build was for when the result names no module: by the
// repository and path the plan's modules were asked at.
func planFailedBuild(ctx context.Context, open *stores, result link.BuildResult) {
entries, err := open.inventory.Catalogued(ctx)
if err != nil {
return
}
for _, e := range entries {
if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path {
planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed)
return
}
}
}
func repositoryMatches(a, b string) bool {
trim := func(s string) string { return strings.ToLower(strings.TrimSuffix(s, ".git")) }
return trim(a) == trim(b) || strings.HasSuffix(trim(a), "/"+trim(b)) || strings.HasSuffix(trim(b), "/"+trim(a))
}
// planStatus is one plan as `status --json` says it.
type planStatus struct {
ID string `json:"id"`
Repository string `json:"repository"`
Commit string `json:"commit"`
State string `json:"state"`
Tier int `json:"tier"`
Tiers int `json:"tiers"`
Waiting string `json:"waiting,omitempty"`
Since time.Time `json:"since"`
Late bool `json:"late"`
}
func planStatuses(plans []inventory.Plan, now time.Time) []planStatus {
out := make([]planStatus, 0, len(plans))
for _, p := range plans {
ps := planStatus{ID: p.ID, Repository: p.Repository, Commit: p.Commit, State: p.State,
Tier: p.Tier, Tiers: len(p.Tiers), Since: p.Updated}
if p.Open() {
ps.Waiting = p.Note
if ps.Waiting == "" {
ps.Waiting = "builds of tier " + fmt.Sprint(p.Tier)
}
ps.Late = now.Sub(p.Updated) > planWaitBound
}
out = append(out, ps)
}
return out
}
// openPlans is the open plans among the recent ones, and how many have waited past the bound.
func openPlans(plans []inventory.Plan) ([]inventory.Plan, int) {
var open []inventory.Plan
late := 0
for _, p := range plans {
if p.Open() {
open = append(open, p)
if time.Since(p.Updated) > planWaitBound {
late++
}
}
}
return open, late
}
// plansCommand says what the last merges produced and where each stands; given an id, one plan
// tier by tier with every module's state.
func plansCommand(ctx context.Context, args []string) error {
set := flag.NewFlagSet("plans", flag.ContinueOnError)
limit := set.Int("n", 10, "how many to show")
whatIf := set.String("what-if", "", "owner/repository: the plan a merge there would produce, saving nothing — with --paths or --modules")
paths := set.String("paths", "", "the files the merge would change, comma-separated, from the repository's root")
modules := set.String("modules", "", "or the modules it would change, comma-separated")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
open, err := openStores(ctx)
if err != nil {
return err
}
defer open.Close()
inv := open.inventory
now := time.Now()
if len(positionals) == 1 {
p, err := inv.PlanByID(ctx, positionals[0])
if err != nil {
return err
}
fmt.Printf("%s — %s\n", p.ID, planLine(p, now))
for i, tier := range p.Tiers {
marker := " "
if i == p.Tier && p.Open() {
marker = ">"
}
fmt.Printf("%s tier %d\n", marker, i)
for _, m := range tier {
s := p.Modules[m]
state := "not yet asked"
if s != nil && s.State != "" {
state = s.State
if s.Commit != "" {
state += " from " + short(s.Commit)
}
if s.Why != "" {
state += ": " + s.Why
}
}
fmt.Printf(" %-22s %s\n", m, state)
}
}
return nil
}
if *whatIf != "" {
return planWhatIf(ctx, inv, *whatIf, splitList(*paths), splitList(*modules))
}
if len(positionals) == 2 && positionals[0] == "stop" {
p, err := inv.PlanByID(ctx, positionals[1])
if err != nil {
return err
}
if !p.Open() {
return fmt.Errorf("%s is already %s", p.ID, p.State)
}
p.State = inventory.PlanFailed
p.Note = "stopped by hand at tier " + fmt.Sprint(p.Tier)
if err := inv.SavePlan(ctx, p); err != nil {
return err
}
fmt.Printf("%s stopped at tier %d of %d; what was asked still builds and registers, nothing further is asked\n",
p.ID, p.Tier, len(p.Tiers))
return nil
}
plans, err := inv.RecentPlans(ctx, *limit)
if err != nil {
return err
}
if len(plans) == 0 {
fmt.Println("no merge has produced a plan yet")
return nil
}
for _, p := range plans {
fmt.Printf("%-28s %s\n", p.ID, planLine(p, now))
}
return nil
}
// planWhatIf is the plan a merge would produce, computed the way the merge handler computes one
// and saved nowhere: the modules the repository's changed files touch (or the modules named), what
// packages their source, everything reachable from them, in tiers. For reading before merging.
func planWhatIf(ctx context.Context, inv *inventory.Inventory, repository string, paths, modules []string) error {
owner, repo, found := strings.Cut(repository, "/")
if !found {
return fmt.Errorf("--what-if takes owner/repository, not %q", repository)
}
m := link.SourceMoved{Owner: owner, Repo: repo, Base: "main", Commit: "what-if", Paths: paths}
entries, err := inv.Catalogued(ctx)
if err != nil {
return err
}
read, err := inv.ReadRepositories(ctx)
if err != nil {
return err
}
var from, packaging []inventory.Entry
named := map[string]bool{}
for _, name := range modules {
named[name] = true
}
for _, e := range entries {
switch {
case named[e.Manifest.Module]:
from = append(from, e)
case len(named) == 0 && sourceIs(e.Source, m):
from = append(from, e)
case readsFrom(read[e.Manifest.Module], m):
packaging = append(packaging, e)
}
}
if len(named) == 0 {
from = whatTheMergeTouched(from, entries, m)
}
moved := append(append([]inventory.Entry{}, from...), packaging...)
if len(moved) == 0 {
fmt.Printf("a merge of %s changing %s would build nothing the mesh holds\n", repository,
orNone(strings.Join(append(paths, modules...), ", ")))
return nil
}
edges, err := inv.Dependencies(ctx)
if err != nil {
return err
}
var names []string
for _, e := range moved {
names = append(names, e.Manifest.Module)
}
p := planOfMerge(m, names, edges)
fmt.Printf("a merge of %s would build %d module(s) in %d tier(s):\n", repository, len(p.Modules), len(p.Tiers))
rolls := map[string]string{}
for i, tier := range p.Tiers {
fmt.Printf(" tier %d\n", i)
for _, name := range tier {
how := "built; its policy records, so nothing is sent"
if u, err := inv.UpgradeOf(ctx, name); err == nil && u.RollOut {
running, _ := inv.Running(ctx, name)
how = "built, then sent to " + orNone(strings.Join(running, ", "))
rolls[name] = how
}
fmt.Printf(" %-22s %s\n", name, how)
}
}
if hasCycle(p.Tiers, edges) {
fmt.Println(" the last tier depends on itself and would be built together, in no order")
}
if len(packaging) > 0 {
var also []string
for _, e := range packaging {
also = append(also, e.Manifest.Module)
}
fmt.Printf(" %s package source from %s, so they are rebuilt without their own source moving\n",
strings.Join(also, ", "), repository)
}
return nil
}
func splitList(s string) []string {
var out []string
for _, part := range strings.Split(s, ",") {
if part = strings.TrimSpace(part); part != "" {
out = append(out, part)
}
}
return out
}
-117
View File
@@ -1,117 +0,0 @@
package main
import (
"testing"
"time"
"github.com/novox/mesh-controller/internal/inventory"
"github.com/novox/mesh-controller/internal/link"
)
// A merge produces a tiered plan (novox/hq ADR 0162): what moved and everything reachable from it,
// sorted so a tier depends only on earlier ones — with the three kinds of dependency told apart.
func TestAMergeIsPlannedInTiersAlongTheThreeKindsOfDependency(t *testing.T) {
edges := []inventory.Edge{
// build dependencies: images on the runtime, a plugin on one of them
{From: "shop", To: "mesh-tools", Kind: inventory.EdgeStandsOn},
{From: "postgres", To: "mesh-tools", Kind: inventory.EdgeStandsOn},
{From: "shop-plugin", To: "shop", Kind: inventory.EdgeDeclared},
// a code dependency: the proxy packages the controller's source — same tier
{From: "route-proxy", To: "mesh-controller", Kind: inventory.EdgePackages},
{From: "builder", To: "mesh-controller", Kind: inventory.EdgePackages},
// runtime dependencies: everything source-built is built by the builder
{From: "shop", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "postgres", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "shop-plugin", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "route-proxy", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "mesh-controller", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "mesh-tools", To: "builder", Kind: inventory.EdgeBuiltBy},
{From: "builder", To: "mesh-tools", Kind: inventory.EdgeStandsOn},
{From: "unrelated", To: "alpine", Kind: inventory.EdgeStandsOn},
}
// The runtime image moved: everything on it, and what is built by what is on it.
set := reachableFrom([]string{"mesh-tools"}, edges)
// What stands on the runtime, and the builder that stands on it; not the controller, which the
// builder merely builds, nor the proxy that packages the controller.
want := []string{"builder", "mesh-tools", "postgres", "shop", "shop-plugin"}
if len(set) != len(want) {
t.Fatalf("reachable from the runtime: %v, want %v", set, want)
}
tiers := tiersOf(set, edges)
pos := map[string]int{}
for i, tier := range tiers {
for _, m := range tier {
pos[m] = i
}
}
if pos["mesh-tools"] != 0 || pos["builder"] != 1 {
t.Fatalf("the runtime then the builder: %v", tiers)
}
if !(pos["shop"] > pos["builder"] && pos["postgres"] > pos["builder"]) {
t.Fatalf("what the builder builds comes after the builder: %v", tiers)
}
if pos["shop-plugin"] <= pos["shop"] {
t.Fatalf("a plugin after what it is declared on: %v", tiers)
}
if hasCycle(tiers, edges) {
t.Fatalf("no cycle here: %v", tiers)
}
// The controller alone moved: the proxy with it, nothing else.
small := reachableFrom([]string{"mesh-controller"}, edges)
if len(small) != 3 {
t.Fatalf("a controller merge rebuilds the controller and what packages it: %v", small)
}
// The builder packages the controller's source (same tier by that edge) and the controller is
// built by the builder (next tier by that one): the builder first, then the controller and the
// proxy together — a code dependency in one tier, a runtime dependency across tiers.
smallTiers := tiersOf(small, edges)
if len(smallTiers) != 2 || smallTiers[0][0] != "builder" || len(smallTiers[1]) != 2 {
t.Fatalf("the builder, then the controller and the proxy together: %v", smallTiers)
}
// The builder alone moved: the builder, and nothing it builds.
if only := reachableFrom([]string{"builder"}, edges); len(only) != 1 {
t.Fatalf("a build machine change rebuilds the build machine alone: %v", only)
}
// Only a runtime dependency gates on deployment, and only when the module rolls out.
p := planOfMerge(link.SourceMoved{Owner: "novox", Repo: "mesh-tools", Commit: "abc"}, []string{"mesh-tools"}, edges)
p.Tier = pos["builder"]
rollsOut := func(m string) bool { return m == "builder" }
if g := gates(p, edges, rollsOut); len(g) != 1 || g[0] != "builder" {
t.Fatalf("the builder gates the tier after it: %v", g)
}
p.Tier = 0
if g := gates(p, edges, rollsOut); len(g) != 0 {
t.Fatalf("the runtime image is a build dependency and gates nothing: %v", g)
}
if g := gates(p, edges, func(string) bool { return false }); len(g) != 0 {
t.Fatalf("a module that only records its upgrade gates nothing: %v", g)
}
}
// A gate is open once every machine running the module has reported after it was built.
func TestAGateOpensWhenTheMachinesHaveReportedSinceTheBuild(t *testing.T) {
built := time.Date(2026, 10, 1, 15, 0, 0, 0, time.UTC)
before, after := built.Add(-time.Minute), built.Add(time.Minute)
reports := []inventory.Reported{{Node: "anchor", At: &after}, {Node: "home-server", At: &before}}
ok, waiting := applied("builder", built, []string{"anchor", "home-server"}, reports)
if ok || len(waiting) != 1 || waiting[0] != "home-server" {
t.Fatalf("one machine has not reported since the build: ok=%v waiting=%v", ok, waiting)
}
if ok, _ := applied("builder", built, []string{"anchor"}, reports); !ok {
t.Fatal("the machine that reported after the build holds the gate open")
}
if ok, _ := applied("builder", built, nil, reports); !ok {
t.Fatal("a module running nowhere gates nothing")
}
}
// A cycle is not lost: what remains is one last tier, and the caller says so.
func TestACycleIsOneLastTierAndSaidSo(t *testing.T) {
edges := []inventory.Edge{{From: "a", To: "b", Kind: inventory.EdgeStandsOn}, {From: "b", To: "a", Kind: inventory.EdgeStandsOn}}
tiers := tiersOf([]string{"a", "b"}, edges)
if len(tiers) != 1 || len(tiers[0]) != 2 || !hasCycle(tiers, edges) {
t.Fatalf("a cycle should be one tier of two, said: %v", tiers)
}
}
+10 -48
View File
@@ -69,24 +69,6 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
return []string{"builds", m}, nil
}
return []string{"builds"}, nil
case "plans":
if r := str("repository"); r != "" {
argv := []string{"plans", "--what-if", r}
if p := str("paths"); p != "" {
argv = append(argv, "--paths", p)
}
if m := str("modules"); m != "" {
argv = append(argv, "--modules", m)
}
return argv, nil
}
if id := str("stop"); id != "" {
return []string{"plans", "stop", id}, nil
}
if id := str("id"); id != "" {
return []string{"plans", id}, nil
}
return []string{"plans"}, nil
case "plan":
if err := need("node"); err != nil {
return nil, err
@@ -97,16 +79,6 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
return nil, err
}
return []string{verb, str("node"), str("module")}, nil
case "pin":
if err := need("node", "provision", "from", "module"); err != nil {
return nil, err
}
return []string{"pin", str("node"), str("provision"), str("from"), str("module")}, nil
case "unpin":
if err := need("node", "provision"); err != nil {
return nil, err
}
return []string{"unpin", str("node"), str("provision")}, nil
case "push":
// Sent and not waited for: the asker reads `status` for what the machine did, which is
// what a person at a shell does too. A tool call that blocked for a push's whole apply would
@@ -130,6 +102,15 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
// the answer the caller needs.
return []string{"rotate"}, nil
case "build":
// Three shapes, as the command has them: a repository, a base every module built on it
// is rebuilt from (`--on`), or everything behind its source (`--behind`). Asked, not
// waited for, the same as a single build.
if on := str("on"); on != "" {
return []string{"build", "--on", on, "--wait", "0"}, nil
}
if b := str("behind"); b != "" && b != "no" && b != "false" {
return []string{"build", "--behind", "--wait", "0"}, nil
}
if err := need("repository"); err != nil {
return nil, err
}
@@ -204,7 +185,7 @@ func seatToolHandlers() (map[string]link.ToolHandler, error) {
}
continue
}
if _, err := argvFor(verb, sampleArguments(v)); err != nil {
if _, err := argvFor(verb, map[string]any{"node": "x", "module": "x", "repository": "x"}); err != nil {
return nil, fmt.Errorf("the %s seat's row declares %q, which this control plane cannot run: %w",
catalogue.ControllerSeatName, verb, err)
}
@@ -243,22 +224,3 @@ func seatTools() map[string]any {
}
return map[string]any{"seats": seats}
}
// sampleArguments is one of every argument a verb's schema requires, so the check at start proves the
// verb runnable rather than that it happens to want the arguments the check guessed.
func sampleArguments(v catalogue.Verb) map[string]any {
sample := map[string]any{"node": "x", "module": "x", "repository": "x"}
switch required := v.Input["required"].(type) {
case []string:
for _, k := range required {
sample[k] = "x"
}
case []any:
for _, k := range required {
if name, ok := k.(string); ok {
sample[name] = "x"
}
}
}
return sample
}
+16
View File
@@ -139,3 +139,19 @@ func TestAJSONVerbsAnswerIsItsStandardOutput(t *testing.T) {
t.Fatalf("stderr and stdout are both what the command said: %s", answer.Output)
}
}
// The build tool has the command's three shapes (ADR 0157's follow-up, 2026-10-01): a repository, a
// base whose dependents are rebuilt, or everything behind its source — each asked, not waited for.
func TestTheBuildToolRebuildsWhatStandsOnABase(t *testing.T) {
argv, _ := argvFor("build", map[string]any{"on": "mesh-tools"})
if strings.Join(argv, " ") != "build --on mesh-tools --wait 0" {
t.Fatalf("a base: %v", argv)
}
argv, _ = argvFor("build", map[string]any{"behind": "yes"})
if strings.Join(argv, " ") != "build --behind --wait 0" {
t.Fatalf("behind: %v", argv)
}
if _, err := argvFor("build", map[string]any{}); err == nil {
t.Fatal("a build naming nothing was accepted")
}
}
+1 -20
View File
@@ -138,18 +138,6 @@ func printStatus(asked answers) error {
len(quiet), strings.Join(said, "\n "))
}
if open, late := openPlans(asked.plans); len(open) > 0 {
fmt.Printf("%d plan(s) open", len(open))
if late > 0 {
fmt.Printf(", %d waiting past %s", late, planWaitBound)
}
fmt.Println(":")
for _, p := range open {
fmt.Printf(" %s\n", planLine(p, time.Now()))
}
fmt.Println()
}
if len(behind) > 0 {
var names []string
for m := range behind {
@@ -301,10 +289,7 @@ func firstLine(s string) string {
// A type of its own rather than a method on the enrolment, because they are unrelated things
// arriving on one queue and an implementation of one should not have to say anything about the
// other.
type builds struct {
inv *inventory.Inventory
open *stores
}
type builds struct{ inv *inventory.Inventory }
// theThreeQuestions reads what anything answering "is the mesh alright" needs.
//
@@ -357,10 +342,6 @@ func theThreeQuestions(ctx context.Context, open *stores) (answers, error) {
if err != nil {
return answers{}, err
}
out.plans, err = inv.RecentPlans(ctx, 5)
if err != nil {
return answers{}, err
}
// And which machines are not running what the mesh would send them. The same question as a
// module being behind its source, one level down: that one says the catalogue is out of date,
-49
View File
@@ -1,49 +0,0 @@
package main
import (
"strings"
"testing"
"github.com/novox/mesh-controller/internal/inventory"
)
// A take is a comparison (novox/hq ADR 0163): the preview puts what runs beside what the module
// declares, and an older image or a differing file refuses unless named.
func TestATakePreviewsTheComparisonAndRefusesWhatIsNotNamed(t *testing.T) {
held := []inventory.Held{
{ID: "forge.server", Module: "forge", Kind: "container", Target: "forge", Facts: map[string]any{
"image": "forge:1.27.3", "image_created": "2026-09-17T10:00:00Z",
"declared_image": "forge:1.22.6", "declared_image_created": "2026-08-20T10:00:00Z", "downgrade": true,
"networks": map[string]any{"predecessor_default": []any{"office", "db"}},
"ports": []any{"3000/tcp>0.0.0.0:3000"}, "declared_ports": []any{"3000:3000"},
}},
{ID: "forge.config", Module: "forge", Kind: "file", Target: "/etc/forge/app.ini", Kept: "/var/lib/mesh/kept/app.ini",
Facts: map[string]any{"differs": true, "difference": []any{"- private scope: local", "+ upstream: public"}}},
{ID: "other.server", Module: "other", Kind: "container", Target: "other"},
}
preview, refusals := comparisonOf(held, "forge", takeOptions{})
for _, want := range []string{"runs forge:1.27.3 (made 2026-09-17)", "declares forge:1.22.6 (made 2026-08-20)", "DOWNGRADE",
"on the network predecessor_default with office, db", "publishes 3000/tcp>0.0.0.0:3000; the module declares 3000:3000",
"- private scope: local", "original kept at /var/lib/mesh/kept/app.ini"} {
if !strings.Contains(preview, want) {
t.Errorf("the preview lacks %q:\n%s", want, preview)
}
}
if strings.Contains(preview, "other") {
t.Errorf("another module's held things are in the preview:\n%s", preview)
}
if len(refusals) != 2 || !strings.Contains(refusals[0], "--downgrade") || !strings.Contains(refusals[1], "--replace /etc/forge/app.ini") {
t.Fatalf("the downgrade and the differing file refuse, each naming its override: %v", refusals)
}
// Named, they pass.
if _, refusals := comparisonOf(held, "forge", takeOptions{Downgrade: true, Replace: map[string]bool{"/etc/forge/app.ini": true}}); len(refusals) != 0 {
t.Fatalf("named differences still refused: %v", refusals)
}
if _, refusals := comparisonOf(held, "forge", takeOptions{Downgrade: true, Replace: map[string]bool{"*": true}}); len(refusals) != 0 {
t.Fatalf("replace * did not cover the file: %v", refusals)
}
// A held thing with no facts yet — a host older than this — refuses nothing and says what it can.
if preview, refusals := comparisonOf(held, "other", takeOptions{}); len(refusals) != 0 || !strings.Contains(preview, "container other") {
t.Fatalf("a factless hold: %q %v", preview, refusals)
}
}
+22 -59
View File
@@ -295,31 +295,17 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
"built from\n", m.Owner, m.Repo, m.Base, m.Commit)
return nil
}
// A merge produces a plan the mesh keeps (novox/hq ADR 0162): what moved and everything that
// depends on it, along the catalogue's one dependency relation, sorted into tiers. The plan is
// written before any build is asked; the first tier is asked; this returns. Outcomes advance it.
edges, err := inv.Dependencies(ctx)
against, err := inv.BuiltAgainst(ctx)
if err != nil {
return notNow(err)
}
var movedNames []string
for _, e := range moved {
movedNames = append(movedNames, e.Manifest.Module)
ordered := orderByBases(moved, against)
names := make([]string, 0, len(ordered))
for _, e := range ordered {
names = append(names, e.Manifest.Module)
}
plan := planOfMerge(m, movedNames, edges)
if hasCycle(plan.Tiers, edges) {
fmt.Printf(" the last tier depends on itself: %s — built together, in no order\n",
strings.Join(plan.Tiers[len(plan.Tiers)-1], ", "))
}
if err := inv.SavePlan(ctx, plan); err != nil {
return notNow(err)
}
var tiers []string
for i, t := range plan.Tiers {
tiers = append(tiers, fmt.Sprintf("%d: %s", i, strings.Join(t, ", ")))
}
fmt.Printf("%s/%s merged into %s (%.8s); plan %s, %d module(s) in %d tier(s)\n %s\n",
m.Owner, m.Repo, m.Base, m.Commit, plan.ID, len(plan.Modules), len(plan.Tiers), strings.Join(tiers, "\n "))
fmt.Printf("%s/%s merged into %s (%.8s); building %s\n",
m.Owner, m.Repo, m.Base, m.Commit, strings.Join(names, ", "))
if len(packaging) > 0 {
var also []string
for _, e := range packaging {
@@ -328,11 +314,22 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
fmt.Printf(" %s package source from it, so they are rebuilt and their own source record "+
"is left where it is\n", strings.Join(also, ", "))
}
if err := askTier(ctx, inv, &plan); err != nil {
return notNow(err)
var failed []string
for _, e := range ordered {
source := buildSource{Repository: e.Source.Repository, Seat: e.Source.Seat}
if err := buildOne(ctx, source, e.Source.Path, e.Source.Ref, 20*time.Minute); err != nil {
fmt.Printf(" %s: %v\n", e.Manifest.Module, err)
failed = append(failed, e.Manifest.Module)
// A base that failed is a reason to stop: what stands on it would be built against
// the old one, and report success (novox/hq 04-ISSUES/131).
if standsOn(ordered, e.Manifest.Module, against) {
fmt.Printf(" stopping: %s is a base of what was still to build\n", e.Manifest.Module)
break
}
}
}
if err := inv.SavePlan(ctx, plan); err != nil {
return notNow(err)
if len(failed) > 0 {
fmt.Printf("%d of %d not built: %s\n", len(failed), len(ordered), strings.Join(failed, ", "))
}
return nil
}
@@ -525,37 +522,3 @@ func isHistory(mergedAt string, seen time.Time) bool {
}
return at.Before(seen)
}
// dependentsOf is every catalogued module that stands on one of the moved modules, directly or
// through another dependent, and is not itself among them — in the catalogue's order, so the
// answer is the same each time. A module standing on nothing that moved is left alone: a merge
// rebuilds what it changed and what is built on top of that, not the catalogue.
func dependentsOf(moved, entries []inventory.Entry, against map[string][]string) []inventory.Entry {
bases := map[string]bool{}
for _, e := range moved {
bases[e.Manifest.Module] = true
}
var out []inventory.Entry
taken := map[string]bool{}
for grew := true; grew; {
grew = false
for _, e := range entries {
name := e.Manifest.Module
if bases[name] || taken[name] {
continue
}
for base := range bases {
if standsOnModule(e, base, against) {
taken[name] = true
out = append(out, e)
grew = true
break
}
}
}
for _, e := range out {
bases[e.Manifest.Module] = true
}
}
return out
}
+11 -28
View File
@@ -716,12 +716,6 @@ func boolByte(b bool) byte {
// routesFrom reads what the mesh wrote and turns it into host → the rules for that host, and
// which of those hosts is a public name — the second is `name`, ACME-eligible; a host reached
// only through `internal-name` never appears there.
//
// **A route may carry either name, or both** (novox/hq ADR 0138). How far an endpoint reaches
// decides which names the mesh composes, so an endpoint that reaches only the private network
// arrives with an `internal-name` and no `name`. That is a whole route, not a malformed one: it is
// served under its internal name and certified by the internal authority. Only a route with
// neither name has nothing to be served under (novox/hq issue 191).
func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
raw, err := os.ReadFile(path)
if err != nil {
@@ -736,18 +730,12 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
public := map[string]bool{}
for _, c := range said.Given {
name, _ := c.Values["name"].(string)
name = strings.TrimSpace(name)
internal, _ := c.Values["internal-name"].(string)
internal = strings.TrimSpace(internal)
if name == "" && internal == "" {
if name == "" {
log.Printf("%s on %s asked for a route and named nothing; skipped", c.From, c.Node)
continue
}
// What the route is called in a log line: its public name when it has one.
called := name
if called == "" {
called = internal
}
host := strings.ToLower(name)
public[host] = true
made := rule{path: asPath(c.Values["path"])}
if p, ok := asWhole(c.Values["priority"]); ok {
@@ -764,7 +752,7 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
if looksLikeACredential(named) {
log.Printf("%s on %s declared route %q with a credential in the declaration rather "+
"than the name of a secret; the whole route is refused (novox/hq ADR 0108)",
c.From, c.Node, called)
c.From, c.Node, name)
continue
}
users, err := usersFrom(named)
@@ -782,7 +770,7 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
port, ok := asPort(c.Values["port"])
if !ok {
log.Printf("%s on %s asked for route %q and gave no usable port; skipped",
c.From, c.Node, called)
c.From, c.Node, name)
continue
}
// Where the mesh says that machine is. Empty means it is this one — a workload beside
@@ -803,7 +791,7 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
}
if scheme != "http" && scheme != "https" {
log.Printf("%s on %s asked for route %q with scheme %q, which is neither http "+
"nor https; skipped", c.From, c.Node, called, scheme)
"nor https; skipped", c.From, c.Node, name, scheme)
continue
}
made.insecure, _ = c.Values["insecure"].(bool)
@@ -814,7 +802,7 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
bytes, whole := asWhole(asked)
if !whole || bytes <= 0 {
log.Printf("%s on %s asked for route %q with a max-request-body of %v, which is "+
"not a whole positive number of bytes; skipped", c.From, c.Node, called, asked)
"not a whole positive number of bytes; skipped", c.From, c.Node, name, asked)
continue
}
made.maxRequestBody = int64(bytes)
@@ -822,20 +810,15 @@ func routesFrom(path string) (map[string][]rule, map[string]bool, error) {
made.target = fmt.Sprintf("%s://%s:%d", scheme, at, port)
}
if name != "" {
host := strings.ToLower(name)
out[host] = append(out[host], made)
public[host] = true
}
out[host] = append(out[host], made)
// The internal-network name, the same rule under a second host — a predecessor proxy
// The internal-network alias, the same rule under a second host — a predecessor proxy
// answered both for one route, as a convenience (reaching a service over the VPN without a
// public TLS round trip), not as an access boundary; composing it here restores exactly
// that, nothing more. Absent whenever the node composed no internal name (novox/hq ADR
// 0056's internalDomain half) — the same "nothing to join a label to" case the public name
// already has. And the only name, when the endpoint reaches no further than the private
// network.
if internal != "" {
// already has.
if internal, _ := c.Values["internal-name"].(string); strings.TrimSpace(internal) != "" {
out[strings.ToLower(internal)] = append(out[strings.ToLower(internal)], made)
}
}
-44
View File
@@ -90,50 +90,6 @@ func TestARouteWithAnInternalNameIsReachableUnderBoth(t *testing.T) {
}
}
// A route whose endpoint reaches only the private network carries an internal name and no public
// one (novox/hq ADR 0138), and is served under that name rather than skipped as naming nothing —
// skipping it left every internal-only module unreachable by name (novox/hq issue 191).
func TestARouteWithOnlyAnInternalNameIsServed(t *testing.T) {
routes, public, err := routesFrom(write(t, `{"given":[
{"from":"app","node":"anchor","at":"anchor.internal",
"values":{"internal-name":"App.Anchor.Internal","port":8443,"scheme":"https","insecure":true}}
]}`))
if err != nil {
t.Fatal(err)
}
if targetOf(routes, "app.anchor.internal") != "https://anchor.internal:8443" {
t.Fatalf("the internal-only route is not served: %v", routes)
}
if len(routes) != 1 {
t.Errorf("an internal-only route made hosts it never named: %v", routes)
}
if len(public) != 0 {
t.Errorf("an internal-only route made a name eligible for a public certificate: %v", public)
}
held := newTable()
held.set(routes, public)
if err := onlyInternalNamesTheMeshSaid(held)(context.Background(), "app.anchor.internal"); err != nil {
t.Errorf("the internal authority refused the internal-only route's name: %v", err)
}
if err := onlyWhatTheMeshSaid(held)(context.Background(), "app.anchor.internal"); err == nil {
t.Error("a public certificate was ordered for an internal-only name")
}
}
// A route with neither name has nothing to be served under, and is still skipped.
func TestARouteWithNeitherNameIsSkipped(t *testing.T) {
routes, public, err := routesFrom(write(t, `{"given":[
{"from":"app","node":"anchor","at":"anchor.internal","values":{"internal-name":" ","port":8080}}
]}`))
if err != nil {
t.Fatal(err)
}
if len(routes) != 0 || len(public) != 0 {
t.Errorf("a route that named nothing was served: %v %v", routes, public)
}
}
// A route with no internal-name composed gets no second host — the ordinary case, unchanged.
func TestARouteWithNoInternalNameGetsNoAlias(t *testing.T) {
routes, _, err := routesFrom(write(t, `{"given":[
+1 -8
View File
@@ -161,15 +161,8 @@ func HolderConsumerFor(node, module string, seat DeclaredSeat) (Consumer, bool)
Queue: "holders",
AckWaitSeconds: 60,
MaxDeliver: 5,
// **One in flight.** A holder works one ask at a time, so the server hands it one at a
// time: with the default of many, every ask behind the one being worked was delivered,
// left unacknowledged for the length of the work, redelivered after the ack wait, and
// after the fifth time dropped — on 2026-10-01 twenty-six of forty-three builds asked in
// two minutes were never built, and the queue read as empty (novox/hq issue 186).
MaxAckPending: 1,
Why: fmt.Sprintf("%s on %s holds %s; it acknowledges after the work is done, so a "+
"crash mid-work redelivers rather than loses; one in flight, so a queue of asks is a "+
"queue and not a race against the ack wait", module, node, seat.Name),
"crash mid-work redelivers rather than loses", module, node, seat.Name),
}, true
}
-13
View File
@@ -153,16 +153,3 @@ func TestANodesDeclarationConsumerIsWhatItsOwnGrantAllows(t *testing.T) {
has(t, perms.Publish, "$JS.ACK.NODES."+c.Name+".>")
has(t, perms.Subscribe, c.Filters[0])
}
// A holder works one ask at a time, so the server hands it one at a time (novox/hq issue 186):
// asks queued behind the one being worked wait in the stream rather than being delivered,
// left to expire and dropped after the fifth redelivery.
func TestAHoldersWorkerTakesOneAskAtATime(t *testing.T) {
c, found := HolderConsumerFor("anchor", "builder", DeclaredSeat{Name: "mesh-build-machine", Accepts: []string{"build"}})
if !found {
t.Fatal("a seat that accepts work has no worker")
}
if c.MaxAckPending != 1 {
t.Fatalf("the worker may have %d asks in flight; one, so a queue is a queue", c.MaxAckPending)
}
}
-1
View File
@@ -144,7 +144,6 @@ func (j *JetStream) EnsureStream(s Stream) error {
MaxMsgsPerSubject: int64(s.MaxMsgsPerSubject),
Description: s.Why,
}
want.AllowDirect = s.Direct
if s.Retention == RetentionLastPerSubject {
// Last-per-subject is a limits stream with one message kept per subject, not a
// retention policy of its own — the state shape, spelled the way the server spells it.
-131
View File
@@ -1,131 +0,0 @@
package broker
import (
"sort"
"strings"
)
// What the mesh issues an assignment to serve and to reach (novox/hq ADR 0160).
//
// A module's code names its tools and its events; **where they land is the mesh's to decide**, and
// it decided it twice — once in the runtime, once here, by one rule compiled into both. Now the
// controller composes a membership for every module on every machine and publishes it to a subject
// only that assignment reads; the runtime serves exactly what the membership says, and the account's
// grant is the same composition read the other way. The shape issued today is the shape the mesh
// already had, so nothing moves when a membership first arrives; only who decides it moves.
// Membership is one assignment's subjects: what this instance of a module on this machine serves,
// and what it may reach.
type Membership struct {
Node string `json:"node"`
Module string `json:"module"`
// Serves is every address a tool of this instance answers on. `{tool}` stands for the tool's
// own name, which the module knows and the mesh does not need to: the mesh issues the address,
// the runtime fills the name. An address with a queue is shared with the module's other
// instances, and the bus hands each call to one of them; an address without is this instance's.
Serves []Served `json:"serves"`
// Seats is every verb of a seat this instance holds, at the subject the seat's callers use.
Seats []SeatServed `json:"seats,omitempty"`
// Emits is where an event of this module lands; `{event}` stands for the event's name.
Emits string `json:"emits"`
// Reaches is each tool this module may call, `<module>.<tool>`, to the subjects that reach it:
// the first is whichever instance answers, when the mesh issued one; the rest name a machine.
Reaches map[string][]string `json:"reaches,omitempty"`
// Tools is where this instance answers what it serves — the runtime's one verb of its own.
Tools string `json:"tools"`
}
// Served is one address a tool is answered on.
type Served struct {
Subject string `json:"subject"`
Queue string `json:"queue,omitempty"`
}
// SeatServed is one verb of a held seat, where its callers ask.
type SeatServed struct {
Seat string `json:"seat"`
Verb string `json:"verb"`
Subject string `json:"subject"`
}
// MembershipSubject is the one address a runtime derives for itself: where its own membership is
// published, from the two names its credential carries. Everything else is in the membership.
func MembershipSubject(node, module string) string {
return "mesh.assignment." + node + "." + module
}
// Placements is where every module runs, for deciding which instance answers for the module.
type Placements struct {
// Nodes is each module's machines.
Nodes map[string][]string
// Interchangeable is each module whose definition says its instances are the same anywhere,
// so the module's plain subject is issued to all of them in one queue.
Interchangeable map[string]bool
}
// AnswersForTheModule says whether an instance of a module on one machine is issued the module's
// plain subject: when it is the only instance, or when the definition says instances are
// interchangeable. A stateful module on two machines gets only its machines' subjects, so a call
// that names none reaches nothing rather than the wrong store.
func (p Placements) AnswersForTheModule(module string) bool {
return len(p.Nodes[module]) <= 1 || p.Interchangeable[module]
}
// MembershipFor composes one assignment's membership from what it declared and where everything
// runs. The subjects are the ones PermissionsFor grants, derived here once more only until the
// grant itself is read from the membership — which is the next step, not this one.
func MembershipFor(node string, d Declared, where Placements) Membership {
own := "mesh.mod." + d.Module
m := Membership{
Node: node, Module: d.Module,
Emits: own + ".event.{event}",
Tools: own + ".tool.tools",
}
// This machine's address always; the module's when this instance answers for the module.
m.Serves = append(m.Serves, Served{Subject: own + ".tool.{tool}." + node})
if where.AnswersForTheModule(d.Module) {
m.Serves = append(m.Serves, Served{Subject: own + ".tool.{tool}", Queue: "serve." + d.Module})
}
for _, s := range d.Holds {
for _, verb := range s.Serves {
m.Seats = append(m.Seats, SeatServed{Seat: s.Name, Verb: verb, Subject: seatToolSubject(s, verb, node)})
}
}
if len(d.Invokes) > 0 {
m.Reaches = map[string][]string{}
for _, t := range d.Invokes {
if t == "*" || strings.HasPrefix(t, "seat:") {
continue // every tool, or a role's: addressed by name, not resolved per instance
}
module, tool, ok := strings.Cut(t, ".")
if !ok {
continue
}
var reach []string
if where.AnswersForTheModule(module) {
reach = append(reach, "mesh.mod."+module+".tool."+tool)
}
nodes := append([]string{}, where.Nodes[module]...)
sort.Strings(nodes)
for _, n := range nodes {
reach = append(reach, "mesh.mod."+module+".tool."+tool+"."+n)
}
m.Reaches[t] = reach
}
}
return m
}
// PlacementsOf reads where everything runs from the records the bus's accounts are composed from.
func PlacementsOf(r Records, interchangeable map[string]bool) Placements {
p := Placements{Nodes: map[string][]string{}, Interchangeable: interchangeable}
for node, declared := range r.Assigned {
for _, d := range declared {
p.Nodes[d.Module] = append(p.Nodes[d.Module], node)
}
}
for _, nodes := range p.Nodes {
sort.Strings(nodes)
}
return p
}
-75
View File
@@ -1,75 +0,0 @@
package broker
import (
"reflect"
"testing"
)
// The mesh issues an assignment's subjects (novox/hq ADR 0160): a module alone on one machine
// answers for the module and for its machine; a stateful module on two machines answers only for
// each machine; one that says its instances are interchangeable answers for the module everywhere;
// a holder serves its seat's verbs; and what a module may reach is resolved the same way.
func TestAMembershipIsIssuedFromWhereEverythingRuns(t *testing.T) {
records := Records{Assigned: map[string][]Declared{
"anchor": {
{Module: "postgres", Serves: []string{"query"}, Holds: []Seat{{Name: "mesh-store", Scope: "mesh", Serves: []string{"databases", "query"}}}},
{Module: "catalog", Invokes: []string{"postgres.query", "search.find"}},
},
"home-server": {
{Module: "postgres"},
{Module: "search"},
{Module: "dashboard", Invokes: []string{"postgres.query"}},
},
"laptop": {{Module: "search"}},
}, Interchangeable: map[string]bool{"search": true}}
where := PlacementsOf(records, records.Interchangeable)
pg := MembershipFor("anchor", records.Assigned["anchor"][0], where)
if !reflect.DeepEqual(pg.Serves, []Served{{Subject: "mesh.mod.postgres.tool.{tool}.anchor"}}) {
t.Fatalf("a stateful module on two machines answers only for its machine: %+v", pg.Serves)
}
if len(pg.Seats) != 2 || pg.Seats[0].Subject != "mesh.seat.mesh-store.tool.databases" {
t.Fatalf("the holder serves the seat's verbs at the seat's subjects: %+v", pg.Seats)
}
if pg.Emits != "mesh.mod.postgres.event.{event}" || pg.Tools != "mesh.mod.postgres.tool.tools" {
t.Fatalf("events and the tools verb: %+v", pg)
}
search := MembershipFor("laptop", records.Assigned["laptop"][0], where)
if !reflect.DeepEqual(search.Serves, []Served{
{Subject: "mesh.mod.search.tool.{tool}.laptop"},
{Subject: "mesh.mod.search.tool.{tool}", Queue: "serve.search"},
}) {
t.Fatalf("an interchangeable module answers for the module in the queue too: %+v", search.Serves)
}
dashboard := MembershipFor("home-server", records.Assigned["home-server"][2], where)
if !reflect.DeepEqual(dashboard.Serves, []Served{
{Subject: "mesh.mod.dashboard.tool.{tool}.home-server"},
{Subject: "mesh.mod.dashboard.tool.{tool}", Queue: "serve.dashboard"},
}) {
t.Fatalf("a module alone on one machine answers for the module: %+v", dashboard.Serves)
}
if !reflect.DeepEqual(dashboard.Reaches["postgres.query"],
[]string{"mesh.mod.postgres.tool.query.anchor", "mesh.mod.postgres.tool.query.home-server"}) {
t.Fatalf("reaching a stateful module names each machine and no plain subject: %v", dashboard.Reaches)
}
catalog := MembershipFor("anchor", records.Assigned["anchor"][1], where)
if !reflect.DeepEqual(catalog.Reaches["search.find"],
[]string{"mesh.mod.search.tool.find", "mesh.mod.search.tool.find.home-server", "mesh.mod.search.tool.find.laptop"}) {
t.Fatalf("reaching an interchangeable module offers the plain subject first: %v", catalog.Reaches)
}
if MembershipSubject("anchor", "postgres") != "mesh.assignment.anchor.postgres" {
t.Fatal("the one subject a runtime derives for itself")
}
}
func TestAnAccountMayReadItsOwnMembershipAndNoOthers(t *testing.T) {
perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "anchor", Module: "postgres", PasswordHash: "x"})
if err != nil {
t.Fatal(err)
}
has(t, perms.Subscribe, "mesh.assignment.anchor.postgres")
has(t, perms.Publish, "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.anchor.postgres")
hasNot(t, perms.Subscribe, "mesh.assignment.>")
}
+2 -8
View File
@@ -173,10 +173,8 @@ func PermissionsFor(p Principal) (Permissions, error) {
switch p.Kind {
case KindController:
// The controller owns the mesh's own traffic and the streams. It is the only writer of
// stream definitions (design 25 §3), so it alone reaches the JetStream API — and it alone
// issues memberships (novox/hq ADR 0160), which it publishes into the assignments stream
// after each push; refused by the server on 2026-10-01 until this line named them.
pub = []string{"mesh.control.>", "mesh.node.>", "mesh.assignment.>", "$JS.API.>"}
// stream definitions (design 25 §3), so it alone reaches the JetStream API.
pub = []string{"mesh.control.>", "mesh.node.>", "$JS.API.>"}
// **And where its consumers deliver.** A push consumer delivers on `_DELIVER.<its name>`,
// and a client bound to it subscribes exactly that; the server refused it for every
// principal the first time one bound a consumer (2026-09-28). Each kind below is granted
@@ -300,10 +298,6 @@ func PermissionsFor(p Principal) (Permissions, error) {
// away — no other principal may subscribe this namespace, and a caller's authority is
// still granted per tool, by name, on the publish side.
sub = append(sub, own+".tool.>")
// Its own membership (ADR 0160): the one subject a runtime derives for itself, read
// directly from the stream and followed live. Nothing else's.
sub = append(sub, MembershipSubject(p.Node, p.Module))
pub = append(pub, "$JS.API.DIRECT.GET."+AssignmentsStream+"."+MembershipSubject(p.Node, p.Module))
// 1b. The tools it calls, if its manifest says it calls any (novox/hq ADR 0152). The same
// grant a person gets and derived the same way, so "what may this module ask" is
-14
View File
@@ -47,14 +47,8 @@ type Stream struct {
// Why is carried into the assertion so an operator reading the server's own state finds the
// reason there, rather than only in a repository they may not have.
Why string
// Direct lets a client read a subject's last message without a consumer, which is how a
// runtime reads its own membership with no JetStream API beyond one request (ADR 0160).
Direct bool
}
// AssignmentsStream holds every assignment's membership, the newest per subject.
const AssignmentsStream = "ASSIGNMENTS"
// MeshStreams is the foundation set, in the order a person reads it.
//
// **CONTROL names its subjects rather than taking `mesh.control.>`**, because heartbeats live
@@ -84,14 +78,6 @@ func MeshStreams() []Stream {
Why: "one declaration per node, always the newest; a node that sees sequence n refuses " +
"n-1 by construction (issue 107)",
},
{
Name: AssignmentsStream,
Subjects: []string{"mesh.assignment.*.*"},
Retention: RetentionLastPerSubject,
Direct: true,
Why: "one membership per assignment, always the newest: what the mesh issued this module " +
"on this machine to serve and to reach (ADR 0160); read directly by the runtime it is for",
},
{
Name: EventsStream,
// A seat's own events ride here too: they are 1:many like any event, and the
+3 -4
View File
@@ -123,10 +123,9 @@ func subjectMatches(filter, subject string) bool {
// Each relationship's retention is the thing that makes it what it is (design 29 §4).
func TestEachStreamCarriesTheRetentionItsShapeNeeds(t *testing.T) {
want := map[string]Retention{
"CONTROL": RetentionWorkQueue,
"NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits,
"ASSIGNMENTS": RetentionLastPerSubject,
"CONTROL": RetentionWorkQueue,
"NODES": RetentionLastPerSubject,
"EVENTS": RetentionLimits,
}
got := map[string]Retention{}
for _, s := range MeshStreams() {
+7 -7
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled
users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused"] }
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "_INBOX.enrol.>", "mesh.control.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.refused"] }
subscribe: { allow: ["$JS.API.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} }
@@ -37,18 +37,18 @@ accounts {
subscribe: { allow: ["_DELIVER.one", "_DELIVER.one.>", "_INBOX.node.one.>", "mesh.node.one.declare"] }
} }
{ user: "one.telegram", password: "$2a$11$tttttttttttttttttttttt", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.one.telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.assignment.one.telegram", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
publish: { allow: ["$JS.ACK.EVENTS.one_telegram.>", "$JS.ACK.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker.>", "$JS.API.CONSUMER.INFO.EVENTS.one_telegram", "$JS.API.CONSUMER.INFO.SEAT_TELEGRAM_SENDER.SEAT_TELEGRAM_SENDER_worker", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.one_telegram", "mesh.seat.telegram-sender.event.delivered", "mesh.seat.telegram-sender.event.failed"] }
subscribe: { allow: ["_DELIVER.SEAT_TELEGRAM_SENDER_worker", "_DELIVER.SEAT_TELEGRAM_SENDER_worker.>", "_INBOX.one.telegram.>", "mesh.mod.telegram.tool.>", "mesh.seat.telegram-sender.accept.send"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "two.audit", password: "$2a$11$aaaaaaaaaaaaaaaaaaaaaa", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.audit"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.assignment.two.audit", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
publish: { allow: ["$JS.ACK.EVENTS.two_audit.>", "$JS.API.CONSUMER.INFO.EVENTS.two_audit", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_audit"] }
subscribe: { allow: ["_INBOX.two.audit.>", "mesh.mod.audit.tool.>", "mesh.mod.shop.event.order.placed"] }
allow_responses: { max: 1, ttl: "1m" }
} }
{ user: "two.shop", password: "$2a$11$ssssssssssssssssssssss", permissions: {
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "$JS.API.DIRECT.GET.ASSIGNMENTS.mesh.assignment.two.shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>", "mesh.assignment.two.shop", "mesh.mod.shop.tool.>"] }
publish: { allow: ["$JS.ACK.EVENTS.two_shop.>", "$JS.API.CONSUMER.INFO.EVENTS.two_shop", "$JS.API.CONSUMER.MSG.NEXT.EVENTS.two_shop", "mesh.mod.shop.event.order.placed", "mesh.seat.telegram-sender.accept.send"] }
subscribe: { allow: ["_INBOX.two.shop.>", "mesh.mod.shop.tool.>"] }
allow_responses: { max: 1, ttl: "1m" }
} }
]
-3
View File
@@ -47,9 +47,6 @@ type Records struct {
Enrolling []string
// People is each person's name against the tools they may invoke, `*` for an administrator.
People map[string][]string
// Interchangeable is each module whose definition says its instances are the same anywhere
// (ADR 0160), which decides whether the module's plain subject is issued to every instance.
Interchangeable map[string]bool
}
// Users is every user the composed file should contain, in the order it will be written.
+5 -5
View File
@@ -25,7 +25,7 @@ func reachable() Node {
func onNetwork(nodes ...string) map[string][]Provider {
out := make([]Provider, 0, len(nodes))
for _, n := range nodes {
out = append(out, Provider{Node: n, At: n + ".internal", Module: "postgres"})
out = append(out, Provider{Node: n, At: n + ".internal"})
}
return map[string][]Provider{"postgres-database": out}
}
@@ -99,7 +99,7 @@ func TestSayingWhichOneSettlesIt(t *testing.T) {
got, err := Resolve(brokeredShelf(), []string{"meshboard"}, reachable(),
World{
Offered: onNetwork("anchor", "archive"),
Pinned: map[string]Chosen{"postgres-database": {Node: "archive", Module: "postgres"}},
Pinned: map[string]string{"postgres-database": "archive"},
})
if err != nil {
t.Fatal(err)
@@ -115,7 +115,7 @@ func TestBeingPointedAtAMachineThatDoesNotProvideItIsRefused(t *testing.T) {
_, err := Resolve(brokeredShelf(), []string{"meshboard"}, reachable(),
World{
Offered: onNetwork("anchor", "archive"),
Pinned: map[string]Chosen{"postgres-database": {Node: "somewhere-else", Module: "postgres"}},
Pinned: map[string]string{"postgres-database": "somewhere-else"},
})
if err == nil {
t.Fatal("a machine was silently given a different database from the one chosen")
@@ -131,12 +131,12 @@ func TestOneProviderDoesNotOverruleAChoice(t *testing.T) {
_, err := Resolve(brokeredShelf(), []string{"meshboard"}, reachable(),
World{
Offered: onNetwork("anchor"),
Pinned: map[string]Chosen{"postgres-database": {Node: "archive", Module: "postgres"}},
Pinned: map[string]string{"postgres-database": "archive"},
})
if err == nil {
t.Fatal("the only database was used although another was chosen")
}
if !strings.Contains(err.Error(), "only anchor/postgres provides it") {
if !strings.Contains(err.Error(), "only anchor provides it") {
t.Fatalf("the refusal does not say what is available: %v", err)
}
}
-100
View File
@@ -1,100 +0,0 @@
package catalogue
import (
"sort"
)
// Chosen is the provider somebody named for a provision: the module, and the node it runs on. Both,
// always (novox/hq #258) — a provision comes from a module, and the same module on two machines is
// two answers, so neither half alone says which. Module is empty only on a record made before this
// was asked, and such a record is honoured exactly as long as it is unambiguous.
type Chosen struct {
Node string
Module string
}
func (c Chosen) String() string {
if c.Module == "" {
return c.Node
}
return c.Node + "/" + c.Module
}
// matches is whether this provider is the one chosen.
func (c Chosen) matches(p Provider) bool {
return p.Node == c.Node && (c.Module == "" || p.Module == c.Module)
}
// among is every offered provider the choice names — one, when the choice is whole.
func (c Chosen) among(where []Provider) []Provider {
var out []Provider
for _, p := range where {
if c.matches(p) {
out = append(out, p)
}
}
return out
}
// nameOf is how a refusal names a provider: the node and the module on it.
func nameOf(p Provider) string {
return Chosen{Node: p.Node, Module: p.Module}.String()
}
// providerNames is every provider named, sorted, for a refusal to list.
func providerNames(where []Provider) []string {
out := make([]string, 0, len(where))
for _, p := range where {
out = append(out, nameOf(p))
}
sort.Strings(out)
return out
}
// providersHere is which modules in this node's own set offer a provision, sorted.
func providersHere(catalogue map[string]Manifest, here func(string) bool, want string) []string {
var out []string
for name, m := range catalogue {
if !here(name) {
continue
}
for _, o := range m.Offers() {
if o == want {
out = append(out, name)
break
}
}
}
sort.Strings(out)
return out
}
// servedByOne is what one provider beside the consumer says a consumer needs to know, or nothing.
//
// Serving is *whether* a need is created at all when the provider is on this same machine (novox/hq
// 04-ISSUES/038's sibling): a need never created is a binding the consumer never gets. The manifest
// alone answers that; the values are settled later, with the node's settings.
func servedByOne(m Manifest, want string) map[string]any {
if _, ok := m.Serves[want]; ok {
return ServedOn(m, want, nil)
}
return nil
}
// sharedByOne is the own secret that provider names as its credential (ADR 0158), or "" when it
// gives each consumer its own.
func sharedByOne(m Manifest, want string) string {
if own, shared := m.SharedCredentialOf(want); shared {
return own
}
return ""
}
func oneOf(list []string, s string) bool {
for _, x := range list {
if x == s {
return true
}
}
return false
}
-21
View File
@@ -107,27 +107,6 @@ func TestCanHoldJudgesClaimScopeAndWhatTheSeatDelivers(t *testing.T) {
if err := CanHold(cannotAnswer, seat); err == nil || !strings.Contains(err.Error(), `does not provide "mesh-bus"`) {
t.Fatalf("a holder that cannot answer for the seat was allowed: %v", err)
}
// A holder's own tools need not be the seat's verbs: the claim may name what it serves for the
// role (ADR 0160), and then only those count — and only the seat's verbs may be named.
promising := seat
promising.Serves = []Verb{{Name: "databases"}, {Name: "query"}}
engine := newBroker()
engine.Tools = []string{"engine_list_databases", "engine_query"}
if err := CanHold(engine, promising); err == nil || !strings.Contains(err.Error(), "does not serve databases, query") {
t.Fatalf("a holder whose tools are not the seat's verbs was allowed without saying what it serves: %v", err)
}
engine.Claims[0].Serves = []string{"databases", "query"}
if err := CanHold(engine, promising); err != nil {
t.Fatalf("a claim naming the seat's verbs was refused: %v", err)
}
engine.Claims[0].Serves = []string{"databases"}
if err := CanHold(engine, promising); err == nil || !strings.Contains(err.Error(), "does not serve query") {
t.Fatalf("a claim naming half the verbs was allowed: %v", err)
}
engine.Claims[0].Serves = []string{"databases", "query", "engine_query"}
if err := CanHold(engine, promising); err == nil || !strings.Contains(err.Error(), "does not promise") {
t.Fatalf("a claim naming a verb the seat never promised was allowed: %v", err)
}
// And the judgement follows the store's row, not a compiled copy.
busSeatDelivering(t, "amqp")
seat, _ = SeatNamed("mesh-broker")
-31
View File
@@ -52,22 +52,6 @@ type Claim struct {
Name string `json:"name"`
// Scope defaults to the node, which is where nearly everything singular is singular.
Scope string `json:"scope,omitempty"`
// Serves names the seat's verbs this module implements for the role, when its own tools are
// not the seat's (novox/hq ADR 0159, 0160): the store's `databases` is not postgres's
// `postgres_list_databases`, and a holder may well serve both. The runtime serves an
// implementation registered under the seat's name on the seat's subjects. Absent, the
// module's own `tools` must list every verb the seat promises, which is how a module named
// like its seat — the catalogue, the records — says they are one and the same.
Serves []string `json:"serves,omitempty"`
}
// ServesFor is what this claim offers a seat's protocol: the verbs it names, else the module's
// own tools.
func (c Claim) ServesFor(m Manifest) []string {
if len(c.Serves) > 0 {
return c.Serves
}
return m.Tools
}
// At is this claim's scope, with the default applied.
@@ -309,12 +293,6 @@ type Manifest struct {
// module claiming a seat answers what that seat's protocol promises (novox/hq ADR 0118).
Tools []string `json:"tools,omitempty"`
// Instances says whether this module's instances are the same anywhere — `interchangeable` —
// so a call that names no machine may be answered by any of them (novox/hq ADR 0160). A fact
// about the software, not about the bus: a stateless web tool says it; a database does not,
// and its instances are then each addressed by machine, never confused for one another.
Instances string `json:"instances,omitempty"`
// Invokes are the tools this module calls, each `<module>.<tool>` or a role's `seat:<seat>.<verb>`,
// or the single entry `*` for every tool on the mesh (novox/hq ADR 0152, ADR 0154).
//
@@ -1195,11 +1173,6 @@ func ParseManifest(raw []byte) (Manifest, error) {
m.Module, r))
}
}
if m.Instances != "" && m.Instances != InstancesInterchangeable {
problems = append(problems, fmt.Sprintf(
"%s says its instances are %q; the one word is %q, for a module that is the same on every machine",
m.Module, m.Instances, InstancesInterchangeable))
}
for _, offer := range m.Provides {
p := offer.Name
if !name.MatchString(p) {
@@ -1987,7 +1960,3 @@ func (o OwnSecrets) Paths() map[string]string {
}
return out
}
// 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"
-115
View File
@@ -1,115 +0,0 @@
package catalogue
import (
"strings"
"testing"
)
// Two modules on one node both provide acme-ca — public-acme (Let's Encrypt) and step-ca (the
// mesh's own authority) on novox — and a route-proxy elsewhere must get the public one (novox/hq
// #258). A pin names the module as well as the node, so that it can say which.
func issuerShelf() map[string]Manifest {
return shelf(
Manifest{Module: "public-acme", Version: "1", Provides: FromAnywhere("acme-ca"),
Serves: map[string]map[string]any{"acme-ca": {"at": "acme-v02.api.letsencrypt.org"}}},
Manifest{Module: "step-ca", Version: "1", Provides: FromAnywhere("acme-ca"),
Serves: map[string]map[string]any{"acme-ca": {"at": "novox.internal"}}},
Manifest{Module: "route-proxy", Version: "1", Requires: []string{"acme-ca"}},
)
}
func twoIssuersOnOneNode() map[string][]Provider {
return map[string][]Provider{"acme-ca": {
{Node: "novox", At: "novox.internal", Module: "public-acme", Serves: map[string]any{"at": "acme-v02.api.letsencrypt.org"}},
{Node: "novox", At: "novox.internal", Module: "step-ca", Serves: map[string]any{"at": "novox.internal"}},
}}
}
func TestTwoProvidersOnOneNodeAreRefusedWithBothNamed(t *testing.T) {
// The refusal must name the pair, because a node alone cannot tell them apart.
_, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode()})
if err == nil {
t.Fatal("two providers on one node were resolved by picking")
}
for _, want := range []string{"novox/public-acme", "novox/step-ca", "<node> <module>"} {
if !strings.Contains(err.Error(), want) {
t.Fatalf("the refusal does not say %q: %v", want, err)
}
}
}
func TestAPinNamesTheModule(t *testing.T) {
got, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode(),
Pinned: map[string]Chosen{"acme-ca": {Node: "novox", Module: "public-acme"}}})
if err != nil {
t.Fatal(err)
}
if len(got.Needs) != 1 || got.Needs[0].From != "novox" || got.Needs[0].Serves["at"] != "acme-v02.api.letsencrypt.org" {
t.Fatalf("the named module was not the one taken: %+v", got.Needs)
}
}
func TestARecordNamingOnlyTheNodeIsRefusedWhenThatNodeAnswersTwice(t *testing.T) {
// A pin from before the module was asked for. It once took the last one listed — a coin flip.
_, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode(), Pinned: map[string]Chosen{"acme-ca": {Node: "novox"}}})
if err == nil {
t.Fatal("a node that answers twice was resolved by picking")
}
if !strings.Contains(err.Error(), "provides it 2 times") || !strings.Contains(err.Error(), "pin workstation acme-ca novox <module>") {
t.Fatalf("the refusal does not ask for the module: %v", err)
}
}
func TestAPinNamingAModuleThatDoesNotProvideItIsRefused(t *testing.T) {
_, err := Resolve(issuerShelf(), []string{"route-proxy"}, reachable(),
World{Offered: twoIssuersOnOneNode(), Pinned: map[string]Chosen{"acme-ca": {Node: "novox", Module: "gitea"}}})
if err == nil || !strings.Contains(err.Error(), "novox/gitea does not provide it") {
t.Fatalf("a module that does not provide it was not refused by name: %v", err)
}
}
func TestTwoProvidersBesideTheConsumerAreRefusedUntilOneIsNamed(t *testing.T) {
// The same ambiguity on the consumer's own machine. This was settled by a map walk — random,
// per plan — which is how novox's own proxy got its issuer.
_, err := Resolve(issuerShelf(), []string{"route-proxy", "public-acme", "step-ca"}, reachable(), World{})
if err == nil {
t.Fatal("two providers beside the consumer were resolved by picking")
}
if !strings.Contains(err.Error(), "workstation provides \"acme-ca\" 2 times") || !strings.Contains(err.Error(), "public-acme, step-ca") {
t.Fatalf("the refusal does not list them: %v", err)
}
}
func TestAPinSettlesTwoProvidersBesideTheConsumer(t *testing.T) {
got, err := Resolve(issuerShelf(), []string{"route-proxy", "public-acme", "step-ca"}, reachable(),
World{Pinned: map[string]Chosen{"acme-ca": {Node: "workstation", Module: "public-acme"}}})
if err != nil {
t.Fatal(err)
}
var found bool
for _, n := range got.Needs {
if n.Name == "acme-ca" {
found = true
if n.Serves["at"] != "acme-v02.api.letsencrypt.org" {
t.Fatalf("the named module was not the one taken: %+v", n)
}
}
}
if !found {
t.Fatalf("no need for acme-ca was created: %+v", got.Needs)
}
}
func TestTheFirstPassDoesNotRefuseTwoProvidersBesideTheConsumer(t *testing.T) {
// The first pass answers only what a node offers. Refused there, the node vanishes from every
// other node's world — and the whole mesh loses its vault for an ambiguity one machine has to
// settle. The second pass is where it is refused, and the test above proves it is.
if _, err := Resolve(issuerShelf(), []string{"route-proxy", "public-acme", "step-ca"}, reachable(),
World{Unchecked: true}); err != nil {
t.Fatalf("the first pass refused what only the second may: %v", err)
}
}
+53 -80
View File
@@ -48,12 +48,11 @@ type World struct {
Holdings []Held
// Offered is what other nodes provide at mesh scope, and everything needed to use it.
Offered map[string][]Provider
// Pinned is which provider this machine was told to get a provision from, by name: a module and
// the node it runs on, both (novox/hq #258). Only consulted when more than one could answer -- a
// choice recorded before it was needed should not start meaning something the day a second
// provider appears, and one recorded and then made unnecessary should not quietly stop applying
// either.
Pinned map[string]Chosen
// Pinned is which node this machine was told to get a provision from, by name. Only consulted
// when more than one node could answer -- a choice recorded before it was needed should not
// start meaning something the day a second provider appears, and one recorded and then made
// unnecessary should not quietly stop applying either.
Pinned map[string]string
// Licences is every provision answered by a **record rather than a node**, by provision name.
//
// novox/hq ADR 0024: a hosted model is on nobody's machine and is reached over the public
@@ -261,10 +260,6 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
// still "choose one" after somebody has chosen one. That makes the remedy useless, and it is
// how this read when first used.
satisfied := map[string]bool{}
// Which modules a person assigned here, hostable. The walk marks a module chosen only when it
// reaches it, and a consumer may be reached before the provider beside it — so the provider of
// something already satisfied is looked for among these as well as among the chosen.
assignedHere := map[string]bool{}
// Everything a person assigned goes in first, except what this machine cannot run. Those are
// choices already made, and a requirement one of them answers is not a choice to put back to
@@ -289,7 +284,6 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
for _, o := range m.Offers() {
satisfied[o] = true
}
assignedHere[a] = true
}
because[a] = "assigned"
queue = append(queue, a)
@@ -317,51 +311,6 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
// commonest arrangement of all — a service and its database on one node — the weakest
// handling, silently.
if satisfied[want] && !isModule(catalogue, want) {
here := func(name string) bool { return chosen[name] || assignedHere[name] }
local := providersHere(catalogue, here, want)
// Which of them it matters to choose between. A plain capability — a shell, a display
// server — asks nothing of whoever answers it, and three shells beside an editor are
// not a choice to put to anybody. One that grants a credential, or serves a fact the
// consumer cannot guess, becomes a binding, and a binding is to one provider.
matter := local
if !brokered[want] {
matter = nil
for _, name := range local {
if _, ok := catalogue[name].Serves[want]; ok {
matter = append(matter, name)
}
}
}
var by Manifest
switch len(matter) {
case 0:
// Nothing to bind to; satisfied by its presence, as it was.
case 1:
by = catalogue[matter[0]]
default:
// Two modules on this machine answer it. Taking whichever a map walk met first
// was the rule until novox/hq #258 — random, per plan — and the same stance as
// across machines applies: ambiguity is refused, never resolved by picking.
//
// **Not in the first pass.** That pass exists only to answer *what does this node
// offer*, and refusing there makes the machine vanish rather than report a problem
// (the sibling case below says why): every other node then loses what this one
// provides — the vault, the identity provider — and refuses for a fault that is
// this node's to settle. The second pass refuses it properly, where it is asked.
if world.Unchecked {
continue
}
c, pinned := world.Pinned[want]
if !pinned || c.Node != node.Name || !oneOf(matter, c.Module) {
reported[want] = true
problems = append(problems, fmt.Sprintf(
"%s provides %q %d times, wanted by %s — say which with `pin %s %s %s <module>`: %s",
node.Name, want, len(matter), because[want], node.Name, want, node.Name,
strings.Join(matter, ", ")))
continue
}
by = catalogue[c.Module]
}
if brokered[want] {
// Answered here, and still a need: the provider is this node.
//
@@ -377,9 +326,9 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
}
needs = append(needs, Needed{
Name: want, From: node.Name, At: at,
Serves: servedByOne(by, want), For: because[want],
SharedOwn: sharedByOne(by, want)})
} else if served := servedByOne(by, want); len(served) > 0 {
Serves: servedHere(catalogue, chosen, want), For: because[want],
SharedOwn: sharedHere(catalogue, chosen, want)})
} else if served := servedHere(catalogue, chosen, want); len(served) > 0 {
// Answered here with no credential to mint, but the provider serves facts the
// consumer cannot guess — a port, a model name — and so still needs a binding.
// **The reachability rule does not apply**: both ends are on this same machine, so
@@ -409,7 +358,11 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
if brokered[want] {
reported[want] = true
where := world.Offered[want]
names := providerNames(where)
names := make([]string, 0, len(where))
for _, p := range where {
names = append(names, p.Node)
}
sort.Strings(names)
take := func(p Provider) {
if node.At != "" && p.At == "" || node.At == "" && p.At != "" || node.At == "" && p.At == "" {
// One of them is not on the private network, so there is no path between
@@ -443,17 +396,17 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
"nothing in this mesh provides %q, wanted by %s %s",
want, because[want], remedy))
case len(where) == 1:
if c, pinned := world.Pinned[want]; pinned && !c.matches(where[0]) {
if chosenNode, pinned := world.Pinned[want]; pinned && chosenNode != where[0].Node {
// One provider, and it is not the one this machine was told to use. Silently
// using the other would be the mesh overruling a choice somebody made.
problems = append(problems, fmt.Sprintf(
"%s was told to get %q from %s, and only %s provides it",
node.Name, want, c, nameOf(where[0])))
node.Name, want, chosenNode, where[0].Node))
break
}
take(where[0])
default:
c, pinned := world.Pinned[want]
chosenNode, pinned := world.Pinned[want]
if !pinned {
// **The seat's holder answers, when a seat delivers this** (novox/hq ADR 0110).
// Not a guess, which ADR 0009 refuses: the choice was made once, mesh-wide, by
@@ -465,32 +418,27 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
break
}
problems = append(problems, fmt.Sprintf(
"%d providers of %q, wanted by %s — say which with `pin %s %s <node> <module>`: %s",
"%d nodes provide %q, wanted by %s — say which with `pin %s %s <node>`: %s",
len(where), want, because[want], node.Name, want,
strings.Join(names, ", ")))
break
}
matching := c.among(where)
switch len(matching) {
case 0:
// Pointed at a provider that does not answer this. Refused rather than
var chosen *Provider
for i, w := range where {
if w.Node == chosenNode {
chosen = &where[i]
}
}
if chosen == nil {
// Pointed at a machine that does not answer this. Refused rather than
// falling back to another: a fallback would quietly move somebody's data to
// a machine they did not choose, which is the whole reason this is asked.
problems = append(problems, fmt.Sprintf(
"%s was told to get %q from %s, and %s does not provide it — these do: %s",
node.Name, want, c, c, strings.Join(names, ", ")))
case 1:
take(matching[0])
default:
// A record naming only the node, from before a pin named the module, and that
// node answers twice. This once took the last one listed (novox/hq #258): a
// coin flip, handed to whoever reads the certificate it chose.
problems = append(problems, fmt.Sprintf(
"%s was told to get %q from %s, and %s provides it %d times — say which with "+
"`pin %s %s %s <module>`: %s",
node.Name, want, c, c.Node, len(matching), node.Name, want, c.Node,
strings.Join(providerNames(matching), ", ")))
node.Name, want, chosenNode, chosenNode, strings.Join(names, ", ")))
break
}
take(*chosen)
}
continue
}
@@ -649,6 +597,31 @@ func Resolve(catalogue map[string]Manifest, assigned []string, node Node, world
// need that is never created is a binding the consumer never gets. It is right about that from the
// manifest alone, which is why walking the catalogue mid-resolution is enough here and is not
// enough for the values.
// sharedHere is the own secret the provider of a provision on this same machine names as its
// credential (ADR 0158), or "" when the provider gives each consumer its own.
func sharedHere(catalogue map[string]Manifest, chosen map[string]bool, want string) string {
for name, m := range catalogue {
if !chosen[name] {
continue
}
if own, shared := m.SharedCredentialOf(want); shared {
return own
}
}
return ""
}
func servedHere(catalogue map[string]Manifest, chosen map[string]bool, want string) map[string]any {
for name, m := range catalogue {
if !chosen[name] {
continue
}
if _, ok := m.Serves[want]; ok {
return ServedOn(m, want, nil)
}
}
return nil
}
// isModule reports whether a name is a module in its own right rather than only something
// modules provide.
+3 -24
View File
@@ -59,26 +59,13 @@ var defaultSeats = []Seat{
{Name: ControllerSeatName, Scope: ScopeMesh, Decision: "novox/hq ADR 0079",
Emits: []string{"applied", "refused", "built-before"},
Serves: ControllerVerbs},
// The store's first verbs (novox/hq ADR 0159): the smallest set that makes the store askable,
// served by whichever module holds the seat with tools of these names.
{Name: "mesh-store", Scope: ScopeMesh, Delivers: "postgres-database", Decision: "novox/hq ADR 0079",
Serves: []Verb{
{Name: "databases", Description: "Every database the store holds, with its on-disk size.",
Input: schema(map[string]string{}, nil)},
{Name: "query", Description: "One read-only statement against one database the store holds.",
Input: schema(map[string]string{"database": "the database to query", "sql": "the read-only statement"},
[]string{"database", "sql"})},
}},
{Name: "mesh-store", Scope: ScopeMesh, Delivers: "postgres-database", Decision: "novox/hq ADR 0079"},
// **Delivers the mesh's own bus, not `amqp`.** Those were the same word until
// ADR 0127 separated them: `amqp` is a backing service a module may require, and this seat is
// the mesh's own transport. ADR 0128 then made that connection something a module requires
// rather than receives ambiently — 23 of the catalogue's modules never speak, and an ambient
// connection would mint a credential for each.
{Name: "mesh-broker", Scope: ScopeMesh, Delivers: "mesh-bus", Decision: "novox/hq ADR 0079"},
// The vault: the controller seals every minted credential with what it provides, which is the
// test for a seat of the mesh's own (novox/hq ADR 0161) — a second provider of `secret` is a
// second claimant, refused by name, rather than a candidate for a pin.
{Name: "mesh-vault", Scope: ScopeMesh, Delivers: "secret", Decision: "novox/hq ADR 0161"},
// Named for its scope since 2026-09-30 (novox/hq ADR 0156); `the-artifact-store` resolves to it as
// an alias on a mesh that predates the rename. It serves artifacts of every kind a build makes —
// images and archives, by digest — which is why the provision is the artifact store and not an
@@ -299,19 +286,11 @@ func CanHold(m Manifest, seat Seat) error {
// **Serving the seat's tools is a condition of holding it** (novox/hq ADR 0132). A holder that
// does not answer what the role promises is every caller's timeout, found at registration and
// at handover instead, naming the verbs rather than the fact that something is missing.
if missing := unservedVerbs(claimed.ServesFor(m), seat.Serves); len(missing) > 0 {
if missing := unservedVerbs(m.Tools, seat.Serves); len(missing) > 0 {
return fmt.Errorf("%s claims %s but does not serve %s, which that seat's protocol promises "+
"(novox/hq ADR 0132) — a holder names every verb its seat declares, under the claim's "+
"serves or among its own tools",
"(novox/hq ADR 0132) — a holder lists every verb its seat declares under tools",
m.Module, seat.Name, strings.Join(missing, ", "))
}
// And nothing the seat does not promise: a verb named here that the protocol lacks is served
// to nobody, which is a typo the holder would otherwise discover as a caller's timeout.
if extra := unpromised(claimed.Serves, seat.Serves); len(extra) > 0 {
return fmt.Errorf("%s claims %s and says it serves %s, which that seat's protocol does not "+
"promise — a claim's serves names the seat's verbs and nothing else",
m.Module, seat.Name, strings.Join(extra, ", "))
}
return nil
}
+4 -4
View File
@@ -208,7 +208,7 @@ func CatalogueProblems(shelf Shelf) []string {
}
// A holder that does not answer what the seat promises is a caller's timeout, found
// at assignment instead.
if missing := unserved(m, c, s); len(missing) > 0 {
if missing := unserved(m, s); len(missing) > 0 {
problems = append(problems, fmt.Sprintf(
"%s claims %s but does not serve %s, which that seat's protocol promises",
module, c.Name, strings.Join(missing, ", ")))
@@ -221,10 +221,10 @@ func CatalogueProblems(shelf Shelf) []string {
// unserved is what a seat's protocol promises and the claimant does not answer. Only the tools
// are checked: `accepts` and `emits` are wired by the runtime from the declaration, while a tool
// is code the module either has or has not written — under the claim's serves, or among its own.
func unserved(m Manifest, c Claim, s SeatDeclaration) []string {
// is code the module either has or has not written.
func unserved(m Manifest, s SeatDeclaration) []string {
has := map[string]bool{}
for _, t := range c.ServesFor(m) {
for _, t := range m.Tools {
has[t] = true
}
var missing []string
-16
View File
@@ -70,22 +70,6 @@ func TestAHolderMustServeWhatItsSeatPromises(t *testing.T) {
}
}
// The claim may say what it serves for the role instead, when the module's own tools are not the
// seat's verbs (ADR 0160).
func TestAClaimMayNameWhatItServesForTheSeat(t *testing.T) {
m := telegram()
m.Tools = []string{"telegram_send"}
m.Claims = append([]Claim(nil), m.Claims...)
for i := range m.Claims {
if m.Claims[i].Name == "telegram-sender" {
m.Claims[i].Serves = []string{"status"}
}
}
if got := problemsFor(t, Shelf{"telegram": m}); strings.Contains(got, "does not serve") {
t.Fatalf("a claim naming the seat's verb was refused: %s", got)
}
}
// A seat with no protocol is a marker: which module is this node's showcase, or its packet filter.
// Most node-scoped seats are markers, so refusing one would refuse the majority of the set.
func TestASeatWithoutAProtocolIsAMarkerNotAMistake(t *testing.T) {
+2 -2
View File
@@ -44,7 +44,7 @@ func TestTheSeatsAreAClosedSetAndEachNamesItsDecision(t *testing.T) {
delivered[s.Delivers] = s.Name
}
}
if len(Seats()) != 15 {
if len(Seats()) != 14 {
t.Errorf("the mesh defines %d seats rather than 14; the set is closed, so a change here is "+
"a decision (novox/hq ADR 0110): %s", len(Seats()), seatNames())
}
@@ -249,7 +249,7 @@ func TestAPinStillWinsOverTheSeat(t *testing.T) {
// A consumer coupled to one provider's contents has said so, and the seat does not overrule it.
got, err := Resolve(registryShelf(), []string{"builder"}, reachable(),
World{Offered: twoRegistries(), Held: giteaHoldsTheSeat(),
Pinned: map[string]Chosen{"npm-package-registry": {Node: "archive", Module: "verdaccio"}}})
Pinned: map[string]string{"npm-package-registry": "archive"}})
if err != nil {
t.Fatal(err)
}
-25
View File
@@ -1,25 +0,0 @@
package catalogue
import "testing"
// The vault's provision is one the controller itself dereferences — every minted credential is
// sealed with it — so it is delivered by a seat of the mesh's own, and a second provider is a second
// claimant refused by name rather than a candidate for a pin (novox/hq ADR 0161, issue 106).
func TestTheVaultsSeatDeliversSecret(t *testing.T) {
seat, known := SeatNamed("mesh-vault")
if !known {
t.Fatal("mesh-vault is not in the mesh's own set")
}
if seat.Scope != ScopeMesh || seat.Delivers != "secret" {
t.Fatalf("mesh-vault is %s-scoped and delivers %q; one per mesh, delivering secret", seat.Scope, seat.Delivers)
}
vault := Manifest{Module: "mesh-vault", Provides: []Offer{{Name: "secret", Scope: ScopeMesh}},
Claims: []Claim{{Name: "mesh-vault", Scope: ScopeMesh}}}
if err := CanHold(vault, seat); err != nil {
t.Fatalf("the vault, claiming its seat and providing secret, was refused: %v", err)
}
another := Manifest{Module: "other-vault", Provides: []Offer{{Name: "secret", Scope: ScopeMesh}}}
if err := CanHold(another, seat); err == nil {
t.Fatal("a provider of secret that does not claim the seat was allowed to hold it")
}
}
+4 -36
View File
@@ -91,31 +91,12 @@ var ControllerVerbs = []Verb{
"module": "one module's name; every module when absent",
"log": "a build's id (as `builds` lists it): print what the build machine said, line by line",
}, nil)},
{Name: "plans", Description: "What the last merges produced and where each stands (novox/hq ADR 0162): " +
"the tiers, the tier a plan is at, what it waits for and since when; one plan whole, given its id.",
Input: schema(map[string]string{
"id": "a plan's id (as `plans` lists them): that plan, tier by tier",
"stop": "a plan's id: stop it — what was asked still builds, nothing further is asked",
"repository": "owner/repository: the plan a merge there would produce, saving nothing (what-if); with paths or modules",
"paths": "with repository: the files the merge would change, comma-separated, from the repository's root",
"modules": "with repository: or the modules it would change, comma-separated",
}, nil)},
{Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it.",
Input: schema(map[string]string{"node": "the machine's name"}, []string{"node"})},
{Name: "assign", Description: "Put a module on a machine. Refused with the mesh's own words when it cannot resolve there.",
Input: schema(map[string]string{"node": "the machine's name", "module": "the module's name"}, []string{"node", "module"})},
{Name: "unassign", Description: "Take a module off a machine.",
Input: schema(map[string]string{"node": "the machine's name", "module": "the module's name"}, []string{"node", "module"})},
{Name: "pin", Description: "Tell a machine which provider answers a provision for it — the module, and the node " +
"it runs on, both. Asked for when more than one could answer; the refusal lists them.",
Input: schema(map[string]string{
"node": "the machine's name",
"provision": "the provision, as the consumer requires it",
"from": "the node the chosen provider runs on",
"module": "the module providing it there",
}, []string{"node", "provision", "from", "module"})},
{Name: "unpin", Description: "Take that choice back, putting the question to the mesh again.",
Input: schema(map[string]string{"node": "the machine's name", "provision": "the provision"}, []string{"node", "provision"})},
{Name: "push", Description: "Send a machine everything it should be — or every machine that is behind, when no machine is named.",
Input: schema(map[string]string{"node": "the machine's name; every machine behind when absent"}, nil)},
{Name: "rotate", Description: "Replace a credential. A pair credential, by provision (and a consuming machine, " +
@@ -132,10 +113,12 @@ var ControllerVerbs = []Verb{
{Name: "build", Description: "Have the build machine build a repository. Answers at once with the build's id: " +
"`builds` with that id follows it line by line, and the module is registered when the outcome comes.",
Input: schema(map[string]string{
"on": "instead of a repository: a module whose artifacts others stand on; every module built on it is rebuilt (the rebuild a changed base needs)",
"behind": "instead of a repository: \"yes\" rebuilds every module the mesh holds older than its source has",
"repository": "the repository's URL, or its path on the forge holding the git seat (owner/name)",
"path": "the module's directory inside it (optional)",
"ref": "the branch, tag or commit to build (optional)",
}, []string{"repository"})},
}, nil)},
}
// schema is a JSON schema for an object of string properties, which is every argument the verbs
@@ -152,22 +135,7 @@ func schema(properties map[string]string, required []string) map[string]any {
return out
}
// unpromised is what a claim says it serves and the seat's protocol never promised.
func unpromised(serves []string, promised []Verb) []string {
has := map[string]bool{}
for _, v := range promised {
has[v.Name] = true
}
var extra []string
for _, s := range serves {
if !has[s] {
extra = append(extra, s)
}
}
return extra
}
// unservedVerbs is what a seat promises and a claimant's offer for it does not answer.
// unservedVerbs is what a seat promises and a claimant's `tools` does not answer.
func unservedVerbs(tools []string, promised []Verb) []string {
has := map[string]bool{}
for _, t := range tools {
+5 -34
View File
@@ -164,18 +164,6 @@ type Held struct {
Since time.Time `json:"since"`
Changed string `json:"changed,omitempty"`
Kept string `json:"kept,omitempty"`
// Facts is what a take compares (novox/hq ADR 0163), as the host reported it: for a found
// container its image and the image's date, the networks and their other members, mounts and
// ports, beside the declared image, ports and volumes, and whether the declared image is the
// older; for a found file whether the declared content differs and how.
Facts map[string]any `json:"facts,omitempty"`
}
// A Stray is a container a machine runs that the mesh neither wrote nor holds (ADR 0163).
type Stray struct {
Kind string `json:"kind"`
Name string `json:"name"`
Detail string `json:"detail,omitempty"`
}
// Reach is one thing reachable on an adopted node: a listening socket or a published port.
@@ -193,8 +181,6 @@ type Adoption struct {
Held []Held
Firewall string
Reachable []Reach
// Strays is what the machine runs that nobody asked for, as last reported (ADR 0163).
Strays []Stray
// At is when it said so; zero when it never has.
At time.Time
}
@@ -203,16 +189,6 @@ type Adoption struct {
// question is the machine as it is now.
func (i *Inventory) RecordAdoption(ctx context.Context, node string, held []Held, firewall string,
reachable []Reach) error {
return i.RecordAdoptionWithStrays(ctx, node, held, firewall, reachable, nil)
}
// RecordAdoptionWithStrays is RecordAdoption with what the machine says strays on it (ADR 0163).
func (i *Inventory) RecordAdoptionWithStrays(ctx context.Context, node string, held []Held, firewall string,
reachable []Reach, strays []Stray) error {
straysRaw, err := json.Marshal(nonNil(strays))
if err != nil {
return err
}
heldRaw, err := json.Marshal(nonNil(held))
if err != nil {
return err
@@ -222,9 +198,9 @@ func (i *Inventory) RecordAdoptionWithStrays(ctx context.Context, node string, h
return err
}
_, err = i.store.Pool().Exec(ctx,
`update node set held = $2, firewall = nullif($3, ''), reachable = $4, strays = $5,
`update node set held = $2, firewall = nullif($3, ''), reachable = $4,
adoption_reported = now(), last_seen = now()
where id = $1`, node, heldRaw, firewall, reachRaw, straysRaw)
where id = $1`, node, heldRaw, firewall, reachRaw)
return err
}
@@ -237,12 +213,12 @@ func nonNil[T any](s []T) []T {
// AdoptionOf is what a node last reported about adoption.
func (i *Inventory) AdoptionOf(ctx context.Context, name string) (Adoption, error) {
var heldRaw, reachRaw, straysRaw []byte
var heldRaw, reachRaw []byte
var firewall *string
var at *time.Time
err := i.store.Pool().QueryRow(ctx,
`select held, firewall, reachable, adoption_reported, strays from node where name = $1`, name).
Scan(&heldRaw, &firewall, &reachRaw, &at, &straysRaw)
`select held, firewall, reachable, adoption_reported from node where name = $1`, name).
Scan(&heldRaw, &firewall, &reachRaw, &at)
if errors.Is(err, pgx.ErrNoRows) {
return Adoption{}, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
@@ -261,11 +237,6 @@ func (i *Inventory) AdoptionOf(ctx context.Context, name string) (Adoption, erro
return Adoption{}, err
}
}
if len(straysRaw) > 0 {
if err := json.Unmarshal(straysRaw, &out.Strays); err != nil {
return Adoption{}, err
}
}
if len(reachRaw) > 0 {
if err := json.Unmarshal(reachRaw, &out.Reachable); err != nil {
return Adoption{}, err
+1 -5
View File
@@ -49,8 +49,7 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
}
}
out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{},
Interchangeable: map[string]bool{}}
out := broker.Records{Assigned: map[string][]broker.Declared{}, People: map[string][]string{}}
for _, n := range nodes {
out.Nodes = append(out.Nodes, n.Name)
modules, err := i.Assigned(ctx, n.Name)
@@ -73,9 +72,6 @@ func (i *Inventory) BusRecords(ctx context.Context) (broker.Records, error) {
"be derived", module, n.Name)
}
out.Assigned[n.Name] = append(out.Assigned[n.Name], declaredFor(m, seats))
if m.Instances == catalogue.InstancesInterchangeable {
out.Interchangeable[m.Module] = true
}
}
}
+16 -21
View File
@@ -841,28 +841,24 @@ func (i *Inventory) SettingsFor(ctx context.Context, nodeName, module string) ([
return layers, rows.Err()
}
// PinProvision records which provider a machine gets a provision from: a module, and the node it
// runs on — both, always (novox/hq #258). A provision comes from a module, and the same module on
// two machines is two answers, so neither half alone says which.
// PinProvision records which node a machine gets a provision from.
//
// Only needed when more than one could answer. Recordable before that, because a mesh with one
// database should not change where an existing machine gets its data the day a second arrives.
func (i *Inventory) PinProvision(ctx context.Context, nodeName, provision, providerNode, module string) error {
if strings.TrimSpace(module) == "" {
return fmt.Errorf("a pin names the module providing %q as well as the node it runs on", provision)
}
// Only needed when more than one node could answer. Recordable before that, because a mesh with
// one database should not change where an existing machine gets its data the day a second
// arrives.
func (i *Inventory) PinProvision(ctx context.Context, nodeName, provision, provider string) error {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return err
}
from, err := i.NodeByName(ctx, providerNode)
from, err := i.NodeByName(ctx, provider)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into provision_pin (node, name, provider, module) values ($1, $2, $3, $4)
on conflict (node, name) do update set provider = excluded.provider, module = excluded.module, pinned_at = now()`,
node.ID, provision, from.ID, module)
`insert into provision_pin (node, name, provider) values ($1, $2, $3)
on conflict (node, name) do update set provider = excluded.provider, pinned_at = now()`,
node.ID, provision, from.ID)
return err
}
@@ -883,28 +879,27 @@ func (i *Inventory) UnpinProvision(ctx context.Context, nodeName, provision stri
return nil
}
// PinsFor is what a node was told about where its provisions come from. A record from before a pin
// named the module carries the node alone; the resolver honours it while it is unambiguous.
func (i *Inventory) PinsFor(ctx context.Context, nodeName string) (map[string]catalogue.Chosen, error) {
// PinsFor is what a node was told about where its provisions come from.
func (i *Inventory) PinsFor(ctx context.Context, nodeName string) (map[string]string, error) {
node, err := i.NodeByName(ctx, nodeName)
if err != nil {
return nil, err
}
rows, err := i.store.Pool().Query(ctx,
`select p.name, n.name, coalesce(p.module, '') from provision_pin p join node n on n.id = p.provider
`select p.name, n.name from provision_pin p join node n on n.id = p.provider
where p.node = $1`, node.ID)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]catalogue.Chosen{}
out := map[string]string{}
for rows.Next() {
var name, provider, module string
if err := rows.Scan(&name, &provider, &module); err != nil {
var name, provider string
if err := rows.Scan(&name, &provider); err != nil {
return nil, err
}
out[name] = catalogue.Chosen{Node: provider, Module: module}
out[name] = provider
}
return out, rows.Err()
}
+4 -4
View File
@@ -385,19 +385,19 @@ func TestAPinSurvivesAndCanBeChanged(t *testing.T) {
t.Fatal(err)
}
}
if err := inv.PinProvision(ctx, "user", "postgres-database", "first", "postgres"); err != nil {
if err := inv.PinProvision(ctx, "user", "postgres-database", "first"); err != nil {
t.Fatal(err)
}
// Changing the answer replaces it rather than adding a second, or a machine would be told to
// use two databases and nothing would say which.
if err := inv.PinProvision(ctx, "user", "postgres-database", "second", "postgres"); err != nil {
if err := inv.PinProvision(ctx, "user", "postgres-database", "second"); err != nil {
t.Fatal(err)
}
pins, err := inv.PinsFor(ctx, "user")
if err != nil {
t.Fatal(err)
}
if len(pins) != 1 || pins["postgres-database"].Node != "second" || pins["postgres-database"].Module != "postgres" {
if len(pins) != 1 || pins["postgres-database"] != "second" {
t.Fatalf("got %v", pins)
}
if err := inv.UnpinProvision(ctx, "user", "postgres-database"); err != nil {
@@ -422,7 +422,7 @@ func TestAPinGoesWhenTheProviderLeavesTheMesh(t *testing.T) {
t.Fatal(err)
}
}
if err := inv.PinProvision(ctx, "consumer", "postgres-database", "provider", "postgres"); err != nil {
if err := inv.PinProvision(ctx, "consumer", "postgres-database", "provider"); err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'provider'`); err != nil {
-124
View File
@@ -1,124 +0,0 @@
package inventory
import (
"context"
"sort"
"strings"
"github.com/novox/mesh-controller/internal/catalogue"
)
// The kinds of edge in the catalogue's one dependency relation (novox/hq ADR 0162).
const (
// EdgeStandsOn: the module's artifact is built on the other's.
EdgeStandsOn = "stands-on"
// EdgePackages: the module's build reads the other's repository.
EdgePackages = "packages"
// EdgeBuiltBy: the module is built by the holder of the build-machine seat.
EdgeBuiltBy = "built-by"
// EdgeDeclared: the manifest's own `build.on`.
EdgeDeclared = "declared"
)
// Edge is one dependency: From depends on To, in the way Kind says.
type Edge struct {
From string `json:"from"`
To string `json:"to"`
Kind string `json:"kind"`
}
// Dependencies is the catalogue's dependency relation, whole: every module the mesh holds, with
// an edge to each module it depends on and the kind of dependency on the edge. One answer, so
// nothing else computes an edge (novox/hq ADR 0162) — the merge handler, `build --on` and the
// overview all read this.
//
// Four sources, one relation: a manifest's `build.on`; the artifacts the latest build was made
// against (an `artifact-store://<module>/…` reference is an edge to that module); the repositories
// the latest build read (an edge to the module whose source that is); and the build machine, which
// every source-built module is built by.
func (i *Inventory) Dependencies(ctx context.Context) ([]Edge, error) {
entries, err := i.Catalogued(ctx)
if err != nil {
return nil, err
}
against, err := i.BuiltAgainst(ctx)
if err != nil {
return nil, err
}
read, err := i.ReadRepositories(ctx)
if err != nil {
return nil, err
}
return dependenciesOf(entries, against, read), nil
}
// dependenciesOf is Dependencies over what was read, so a test can hand it a catalogue.
func dependenciesOf(entries []Entry, against map[string][]string, read map[string][]ReadRepository) []Edge {
known := map[string]bool{}
byRepository := map[string][]string{}
var builders []string
for _, e := range entries {
name := e.Manifest.Module
known[name] = true
if r := repositoryKey(e.Source.Repository); r != "" {
byRepository[r] = append(byRepository[r], name)
}
if e.Manifest.ClaimsSeat("mesh-build-machine") {
builders = append(builders, name)
}
}
seen := map[Edge]bool{}
var out []Edge
add := func(from, to, kind string) {
if from == to || !known[to] {
return
}
e := Edge{From: from, To: to, Kind: kind}
if !seen[e] {
seen[e] = true
out = append(out, e)
}
}
for _, e := range entries {
name := e.Manifest.Module
if e.Manifest.Build != nil {
for _, on := range e.Manifest.Build.On {
if on.Module != "" {
add(name, on.Module, EdgeDeclared)
}
}
}
for _, ref := range against[name] {
if rest, ok := strings.CutPrefix(ref, catalogue.ArtifactStoreScheme); ok {
if base, _, found := strings.Cut(rest, "/"); found {
add(name, base, EdgeStandsOn)
}
}
}
for _, r := range read[name] {
for _, other := range byRepository[repositoryKey(r.Repository)] {
add(name, other, EdgePackages)
}
}
if e.Source.Repository != "" {
for _, b := range builders {
add(name, b, EdgeBuiltBy)
}
}
}
sort.Slice(out, func(a, b int) bool {
if out[a].From != out[b].From {
return out[a].From < out[b].From
}
if out[a].To != out[b].To {
return out[a].To < out[b].To
}
return out[a].Kind < out[b].Kind
})
return out
}
// repositoryKey is a repository as compared: lower-cased, without a trailing `.git`.
func repositoryKey(repository string) string {
return strings.ToLower(strings.TrimSuffix(strings.TrimSpace(repository), ".git"))
}
-64
View File
@@ -1,64 +0,0 @@
package inventory
import (
"testing"
"github.com/novox/mesh-controller/internal/catalogue"
)
// The catalogue's one dependency relation (novox/hq ADR 0162): four kinds of edge from one call.
func TestDependenciesAreOneRelationWithTheirKinds(t *testing.T) {
entry := func(name, repository string) Entry {
return Entry{Manifest: catalogue.Manifest{Module: name}, Source: Source{Repository: repository}}
}
builder := entry("builder", "http://forge/novox/mesh-catalog.git")
builder.Manifest.Claims = []catalogue.Claim{{Name: "mesh-build-machine", Scope: catalogue.ScopeMesh}}
plugin := entry("shop-plugin", "http://forge/novox/mesh-catalog.git")
plugin.Manifest.Build = &catalogue.Build{On: []catalogue.BuildsOn{{Arg: "BASE", Module: "shop"}}}
entries := []Entry{
entry("mesh-tools", "http://forge/novox/mesh-tools.git"),
entry("shop", "http://forge/novox/mesh-catalog.git"),
plugin,
builder,
entry("mesh-controller", "http://forge/novox/mesh-controller.git"),
entry("route-proxy", "http://forge/novox/mesh-catalog.git"),
{Manifest: catalogue.Manifest{Module: "hand-made"}},
}
against := map[string][]string{
"shop": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
"builder": {catalogue.ArtifactStoreScheme + "mesh-tools/runtime@sha256:a"},
}
read := map[string][]ReadRepository{
"route-proxy": {{Repository: "http://forge/novox/mesh-controller.git", Ref: "main"}},
}
got := dependenciesOf(entries, against, read)
has := func(from, to, kind string) bool {
for _, e := range got {
if e == (Edge{From: from, To: to, Kind: kind}) {
return true
}
}
return false
}
for _, want := range []Edge{
{"shop", "mesh-tools", EdgeStandsOn},
{"builder", "mesh-tools", EdgeStandsOn},
{"shop-plugin", "shop", EdgeDeclared},
{"route-proxy", "mesh-controller", EdgePackages},
{"shop", "builder", EdgeBuiltBy},
{"mesh-controller", "builder", EdgeBuiltBy},
{"mesh-tools", "builder", EdgeBuiltBy},
} {
if !has(want.From, want.To, want.Kind) {
t.Errorf("missing %+v in %+v", want, got)
}
}
if has("builder", "builder", EdgeBuiltBy) {
t.Error("the builder is not built by itself")
}
for _, e := range got {
if e.From == "hand-made" {
t.Errorf("a module with no source depends on nothing: %+v", e)
}
}
}
@@ -1,12 +0,0 @@
-- A pin names the module as well as the node (novox/hq #258).
--
-- 0008 said "not a module: the same module on two machines is two answers, and which machine is the
-- whole question". Half right. Two modules on one machine can both answer a provision — public-acme
-- and step-ca both offer acme-ca on novox — and then which *module* is the whole question, and a
-- node alone cannot ask it. The resolver, given a node that answered twice, took the last one listed.
--
-- A provider is a (node, module) pair (design 23), and a pin names the pair. Nullable, so a record
-- made before this was asked keeps meaning what it meant: honoured while that node answers once,
-- refused with the module asked for when it answers twice.
alter table provision_pin add column module text;
@@ -1,19 +0,0 @@
-- The records already made are completed where the mesh can tell: a pin naming a node on which
-- exactly one assigned module offers the provision (or is the module itself, for a requirement that
-- names a module) gets that module. A node that answers twice is left to say which — the resolver
-- refuses it with the module asked for, rather than this guessing on its behalf.
update provision_pin p
set module = sub.module
from (
select p2.node, p2.name, min(a.module) as module, count(distinct a.module) as answers
from provision_pin p2
join assignment a on a.node = p2.provider
join module m on m.name = a.module
where p2.module is null
and (a.module = p2.name
or exists (select 1
from jsonb_array_elements(coalesce(m.manifest -> 'provides', '[]'::jsonb)) e
where (case when jsonb_typeof(e) = 'string' then e #>> '{}' else e ->> 'name' end) = p2.name))
group by p2.node, p2.name
) sub
where sub.node = p.node and sub.name = p.name and sub.answers = 1;
@@ -1,19 +0,0 @@
-- A merge produces a tiered plan the mesh keeps (novox/hq ADR 0162): what the merge changed and
-- everything standing on it, sorted into tiers, each module's state, and the tier the plan is at.
-- Kept so a controller replaced mid-plan resumes it, and so `status` can say what a merge still
-- waits for.
create table release_plan (
id text primary key,
repository text not null,
commit_hash text not null,
created timestamptz not null default now(),
updated timestamptz not null default now(),
-- building: a tier's builds are asked; rolling: the tier is built and the machines are applying
-- what a later tier needs running; done; failed.
state text not null,
tier int not null default 0,
tiers jsonb not null,
modules jsonb not null,
note text not null default ''
);
create index release_plan_open on release_plan (created) where state in ('building', 'rolling');
@@ -1,3 +0,0 @@
-- What runs on a machine that the mesh neither wrote nor holds, as the host reports it with every
-- apply (novox/hq ADR 0163): a container left behind by a cutover is seen the day it is left.
alter table node add column strays jsonb;
+2 -5
View File
@@ -42,11 +42,8 @@ func TestAPersonMayCallToolsAndNothingElse(t *testing.T) {
if err != nil {
t.Fatal(err)
}
// The one tool, both ways it is addressed (novox/hq ADR 0159): to whichever instance
// answers, and to the instance on one machine. Nothing else.
if len(perms.Publish) != 2 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" ||
perms.Publish[1] != "mesh.mod.mesh-catalog.tool.catalog_tools.*" {
t.Errorf("ada may publish %v, which should be the one tool, both ways addressed, and nothing else", perms.Publish)
if len(perms.Publish) != 1 || perms.Publish[0] != "mesh.mod.mesh-catalog.tool.catalog_tools" {
t.Errorf("ada may publish %v, which should be the one tool and nothing else", perms.Publish)
}
for _, s := range perms.Publish {
if strings.HasPrefix(s, "mesh.control") || strings.HasPrefix(s, "mesh.node") ||
-86
View File
@@ -1,86 +0,0 @@
package inventory
import (
"context"
"os"
"testing"
)
// A pin made before it named the module (migration 0051, novox/hq #258): completed where the node it
// names answers once, left for a person where it answers twice.
func legacyPin(t *testing.T, inv *Inventory, node, provision, provider string) {
t.Helper()
_, err := inv.store.Pool().Exec(context.Background(),
`insert into provision_pin (node, name, provider)
select u.id, $2, p.id from node u, node p where u.name = $1 and p.name = $3`,
node, provision, provider)
if err != nil {
t.Fatal(err)
}
}
func completeEarlierPins(t *testing.T, inv *Inventory) {
t.Helper()
sql, err := os.ReadFile("migrations/0051-a-pin-made-before-is-completed.sql")
if err != nil {
t.Fatal(err)
}
if _, err := inv.store.Pool().Exec(context.Background(), string(sql)); err != nil {
t.Fatal(err)
}
}
func TestAPinMadeBeforeIsCompletedWhenTheNodeAnswersOnce(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
for _, n := range []string{"user", "provider"} {
if _, err := inv.AddNode(ctx, n); err != nil {
t.Fatal(err)
}
}
for _, m := range []string{"postgres", "redis"} {
if err := inv.RegisterModule(ctx, manifest(m, []string{m + "-database"}, nil), Source{}); err != nil {
t.Fatal(err)
}
if _, err := inv.Assign(ctx, "provider", m); err != nil {
t.Fatal(err)
}
}
legacyPin(t, inv, "user", "postgres-database", "provider")
completeEarlierPins(t, inv)
pins, err := inv.PinsFor(ctx, "user")
if err != nil {
t.Fatal(err)
}
if got := pins["postgres-database"]; got.Node != "provider" || got.Module != "postgres" {
t.Fatalf("the record was not completed with the one module that answers: %+v", got)
}
}
func TestAPinMadeBeforeIsLeftOpenWhenTheNodeAnswersTwice(t *testing.T) {
inv := fresh(t)
ctx := context.Background()
for _, n := range []string{"user", "provider"} {
if _, err := inv.AddNode(ctx, n); err != nil {
t.Fatal(err)
}
}
for _, m := range []string{"public-acme", "step-ca"} {
if err := inv.RegisterModule(ctx, manifest(m, []string{"acme-ca"}, nil), Source{}); err != nil {
t.Fatal(err)
}
if _, err := inv.Assign(ctx, "provider", m); err != nil {
t.Fatal(err)
}
}
legacyPin(t, inv, "user", "acme-ca", "provider")
completeEarlierPins(t, inv)
pins, err := inv.PinsFor(ctx, "user")
if err != nil {
t.Fatal(err)
}
if got := pins["acme-ca"]; got.Node != "provider" || got.Module != "" {
t.Fatalf("a node that answers twice was guessed for: %+v", got)
}
}
-127
View File
@@ -1,127 +0,0 @@
package inventory
import (
"context"
"encoding/json"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
)
// A Plan is what a merge produces (novox/hq ADR 0162): the modules it changed and everything
// standing on them, sorted into tiers, each module's state, and the tier the plan is at. Kept in
// the store so a controller replaced mid-plan resumes it, and so `status` can say what a merge
// still waits for.
type Plan struct {
ID string `json:"id"`
Repository string `json:"repository"`
Commit string `json:"commit"`
Created time.Time `json:"created"`
Updated time.Time `json:"updated"`
State string `json:"state"`
Tier int `json:"tier"`
Tiers [][]string `json:"tiers"`
Modules map[string]*PlanModule `json:"modules"`
Note string `json:"note,omitempty"`
}
// PlanModule is one module's state within a plan.
type PlanModule struct {
// State: asked, built, failed; empty for a module whose tier has not been asked yet.
State string `json:"state,omitempty"`
AskedAt *time.Time `json:"asked_at,omitempty"`
BuiltAt *time.Time `json:"built_at,omitempty"`
// SentAt is when the plan sent the machines running this module its new build, because a
// later tier is built by it (ADR 0163's gate): the reports that open the gate are the ones
// after this.
SentAt *time.Time `json:"sent_at,omitempty"`
Commit string `json:"commit,omitempty"`
Why string `json:"why,omitempty"`
}
// The states a plan passes through.
const (
PlanBuilding = "building"
PlanRolling = "rolling"
PlanDone = "done"
PlanFailed = "failed"
)
// Open says whether the plan is still being worked.
func (p Plan) Open() bool { return p.State == PlanBuilding || p.State == PlanRolling }
// SavePlan writes a plan, new or changed, whole: the plan is small and read as one thing.
func (i *Inventory) SavePlan(ctx context.Context, p Plan) error {
tiers, err := json.Marshal(p.Tiers)
if err != nil {
return err
}
modules, err := json.Marshal(p.Modules)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`insert into release_plan (id, repository, commit_hash, created, updated, state, tier, tiers, modules, note)
values ($1, $2, $3, $4, now(), $5, $6, $7, $8, $9)
on conflict (id) do update set updated = now(), state = excluded.state, tier = excluded.tier,
tiers = excluded.tiers, modules = excluded.modules, note = excluded.note`,
p.ID, p.Repository, p.Commit, p.Created, p.State, p.Tier, tiers, modules, p.Note)
return err
}
// OpenPlans is every plan still being worked, oldest first.
func (i *Inventory) OpenPlans(ctx context.Context) ([]Plan, error) {
return i.plans(ctx, `where state in ('building', 'rolling') order by created`)
}
// RecentPlans is the last few plans, newest first, open or not — what the overview shows.
func (i *Inventory) RecentPlans(ctx context.Context, limit int) ([]Plan, error) {
return i.plans(ctx, fmt.Sprintf(`order by created desc limit %d`, limit))
}
// PlanByID is one plan.
func (i *Inventory) PlanByID(ctx context.Context, id string) (Plan, error) {
plans, err := i.plans(ctx, `where id = '`+id+`'`)
if err != nil {
return Plan{}, err
}
if len(plans) == 0 {
return Plan{}, fmt.Errorf("no plan %s", id)
}
return plans[0], nil
}
func (i *Inventory) plans(ctx context.Context, tail string) ([]Plan, error) {
rows, err := i.store.Pool().Query(ctx,
`select id, repository, commit_hash, created, updated, state, tier, tiers, modules, note
from release_plan `+tail)
if err != nil {
return nil, err
}
defer rows.Close()
var out []Plan
for rows.Next() {
var p Plan
var tiers, modules []byte
if err := rows.Scan(&p.ID, &p.Repository, &p.Commit, &p.Created, &p.Updated, &p.State,
&p.Tier, &tiers, &modules, &p.Note); err != nil {
return nil, err
}
if err := json.Unmarshal(tiers, &p.Tiers); err != nil {
return nil, err
}
if err := json.Unmarshal(modules, &p.Modules); err != nil {
return nil, err
}
if p.Modules == nil {
p.Modules = map[string]*PlanModule{}
}
out = append(out, p)
}
if errors.Is(rows.Err(), pgx.ErrNoRows) {
return nil, nil
}
return out, rows.Err()
}
-48
View File
@@ -1,48 +0,0 @@
package inventory
import (
"testing"
"time"
)
// A plan is a record the mesh keeps and resumes (novox/hq ADR 0162): written whole, read back open,
// advanced, and gone from the open ones when done.
func TestAPlanIsKeptAdvancedAndResumedFromTheStore(t *testing.T) {
inv := ForTest(t)
ctx := t.Context()
p := Plan{ID: "plan-1", Repository: "novox/mesh-tools", Commit: "abc", Created: time.Now().UTC(),
State: PlanBuilding, Tiers: [][]string{{"mesh-tools"}, {"builder"}, {"shop"}},
Modules: map[string]*PlanModule{"mesh-tools": {}, "builder": {}, "shop": {}}}
if err := inv.SavePlan(ctx, p); err != nil {
t.Fatal(err)
}
open, err := inv.OpenPlans(ctx)
if err != nil || len(open) != 1 || open[0].ID != "plan-1" || len(open[0].Tiers) != 3 {
t.Fatalf("the plan was not kept whole: %v %+v", err, open)
}
// Another controller picks it up where it was left: a tier advanced and a module built.
now := time.Now().UTC()
resumed := open[0]
resumed.Tier = 1
resumed.Modules["mesh-tools"].State = "built"
resumed.Modules["mesh-tools"].BuiltAt = &now
resumed.State = PlanRolling
resumed.Note = "tier 0 built; waiting for builder on anchor to be applied"
if err := inv.SavePlan(ctx, resumed); err != nil {
t.Fatal(err)
}
again, err := inv.PlanByID(ctx, "plan-1")
if err != nil || again.Tier != 1 || again.Modules["mesh-tools"].State != "built" || again.State != PlanRolling {
t.Fatalf("the advanced plan did not come back as left: %v %+v", err, again)
}
again.State = PlanDone
if err := inv.SavePlan(ctx, again); err != nil {
t.Fatal(err)
}
if open, _ = inv.OpenPlans(ctx); len(open) != 0 {
t.Fatalf("a done plan is not open: %+v", open)
}
if recent, _ := inv.RecentPlans(ctx, 5); len(recent) != 1 || recent[0].State != PlanDone {
t.Fatalf("a done plan is still among the recent ones: %+v", recent)
}
}
+2 -2
View File
@@ -398,7 +398,7 @@ func (i *Inventory) AcceptSecretForModule(ctx context.Context, node, module, nam
return err
}
if _, own := m.OwnSecrets[name]; !own {
return fmt.Errorf("%s does not declare %q as an own secret; %s — a secret it requires from a provider is accepted with `--provider <node> [--local <name>]`, the value the running service already uses (novox/hq ADR 0163)", module, name, declaresOwn(m))
return fmt.Errorf("%s does not declare %q as an own secret; %s", module, name, declaresOwn(m))
}
key, err := i.SealingKeyOf(ctx, node)
if err != nil {
@@ -546,7 +546,7 @@ func (i *Inventory) RotateModuleSecret(ctx context.Context, node, module, name s
}
own, declared := m.OwnSecrets[name]
if !declared {
return fmt.Errorf("%s does not declare %q as an own secret; %s — a secret it requires from a provider is accepted with `--provider <node> [--local <name>]`, the value the running service already uses (novox/hq ADR 0163)", module, name, declaresOwn(m))
return fmt.Errorf("%s does not declare %q as an own secret; %s", module, name, declaresOwn(m))
}
switch own.Taken {
case catalogue.TakenAtStart:
-6
View File
@@ -173,13 +173,7 @@ func (m *natsMachine) Take(ctx context.Context, do func(context.Context, Build))
_ = msg.Term()
continue
}
// A build outlives the acknowledgement window many times over; said while it runs,
// as the controller says it for its own long handlers, so the server neither hands
// the ask to a second machine nor counts the wait against its deliveries.
working := make(chan struct{})
go stillWorking(msg, working)
do(ctx, &natsBuild{request: request, msg: msg, on: m.on, js: m.js})
close(working)
}
}
}
-19
View File
@@ -4,7 +4,6 @@ import (
"context"
"errors"
"fmt"
"github.com/novox/mesh-controller/internal/broker"
"time"
"github.com/nats-io/nats.go"
@@ -146,24 +145,6 @@ func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body
return nil
}
// PublishMembership issues one assignment what it serves and reaches (novox/hq ADR 0160), last per
// subject, so the runtime that connects later reads the current one and one that is running follows.
// MembershipWait bounds how long issuing one membership may take. A publish the server refuses is
// never acknowledged, and a stream publish waits for its acknowledgement for as long as its
// context lives: on 2026-10-01 the daemon's own context was that long, and one refused membership
// held the controller's receive loop for good (novox/hq issue 185).
const MembershipWait = 10 * time.Second
func (b OverNATS) PublishMembership(ctx context.Context, node, module string, body []byte) error {
ctx, cancel := context.WithTimeout(ctx, MembershipWait)
defer cancel()
_, err := b.JS.Publish(broker.MembershipSubject(node, module), body, nats.Context(ctx))
if err != nil {
return fmt.Errorf("issuing %s on %s its membership: %w", module, node, err)
}
return nil
}
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
if err != nil {
+3 -14
View File
@@ -286,22 +286,18 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err
// schedule, when what it holds changes, not only after an apply — and never cleared by a
// report that carries none, which is every bare word that the node is there. An adopted node
// always names its firewall, so a report from one replaces all three, emptied held included.
if len(report.Held) > 0 || report.Firewall != "" || len(report.Reachable) > 0 || len(report.Strays) > 0 {
if len(report.Held) > 0 || report.Firewall != "" || len(report.Reachable) > 0 {
held := make([]inventory.Held, 0, len(report.Held))
for _, h := range report.Held {
held = append(held, inventory.Held{ID: h.ID, Module: h.Module, Kind: h.Kind,
Target: h.Target, Since: h.Since, Changed: h.Changed, Kept: h.Kept, Facts: h.Facts})
}
strays := make([]inventory.Stray, 0, len(report.Strays))
for _, s := range report.Strays {
strays = append(strays, inventory.Stray{Kind: s.Kind, Name: s.Name, Detail: s.Detail})
Target: h.Target, Since: h.Since, Changed: h.Changed, Kept: h.Kept})
}
reachable := make([]inventory.Reach, 0, len(report.Reachable))
for _, r := range report.Reachable {
reachable = append(reachable, inventory.Reach{Protocol: r.Protocol, Address: r.Address,
Port: r.Port, By: r.By, Published: r.Published, ContainerPort: r.ContainerPort})
}
if err := e.Inventory.RecordAdoptionWithStrays(ctx, node.ID, held, report.Firewall, reachable, strays); err != nil {
if err := e.Inventory.RecordAdoption(ctx, node.ID, held, report.Firewall, reachable); err != nil {
return false, err
}
}
@@ -324,13 +320,6 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err err
return false, err
}
}
if len(report.Profile) > 0 {
// The latest wins, as at enrolment: a capability the machine lost is one the plan must
// stop counting on (novox/hq ADR 0161).
if err := e.Inventory.RecordProfile(ctx, node.ID, report.Profile); err != nil {
return false, err
}
}
// What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it.
if report.Tunnel != nil {
if err := e.Inventory.RecordCarriedTunnel(ctx, node.ID, inventory.Carried{
-45
View File
@@ -1,45 +0,0 @@
package link_test
import (
"testing"
"github.com/novox/mesh-controller/internal/link"
)
// A report may carry the machine's profile, detected again by the apply that reports, and the latest
// replaces what enrolment recorded (novox/hq ADR 0161): a machine that switched its network manager
// is a machine whose uplink holder lacks a capability at its next push, not at its next enrolment.
func TestAReportsProfileReplacesTheEnrolledOne(t *testing.T) {
e, _, _ := anEnrolledHub(t)
ctx := t.Context()
first := map[string]any{"capabilities": []any{map[string]any{"name": "uplink-networkmanager", "present": true}}}
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Profile: first}); err != nil {
t.Fatal(err)
}
got, err := e.Inventory.Profile(ctx, "anchor")
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Name != "uplink-networkmanager" || !got[0].Present {
t.Fatalf("the report's profile was not kept: %+v", got)
}
// The machine switched managers; the next report says so and the old fact is gone.
second := map[string]any{"capabilities": []any{map[string]any{"name": "uplink-systemd-networkd", "present": true}}}
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Profile: second}); err != nil {
t.Fatal(err)
}
got, err = e.Inventory.Profile(ctx, "anchor")
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Name != "uplink-systemd-networkd" {
t.Fatalf("the latest profile did not replace the earlier one: %+v", got)
}
// A report with no profile leaves the last one standing.
if _, err := e.Heard(ctx, link.Report{Node: "anchor", Host: "1"}); err != nil {
t.Fatal(err)
}
if got, _ = e.Inventory.Profile(ctx, "anchor"); len(got) != 1 {
t.Fatalf("a report without a profile erased it: %+v", got)
}
}
-19
View File
@@ -193,15 +193,6 @@ type Report struct {
// refuses it whole — which is right, and makes every new field a flag day that the mesh could
// not see coming.
Host string `json:"host,omitempty"`
// Strays is what runs on the machine that the mesh neither wrote nor holds (ADR 0163).
Strays []Stray `json:"strays,omitempty"`
// Profile is what the machine can do, detected again by this apply (novox/hq ADR 0161): the
// same shape enrolment sends, so a machine that gained or lost a capability — switched its
// network manager — is known at its next push and not at its next enrolment. Absent from a host
// older than this, and then the enrolment's profile stands.
Profile map[string]any `json:"profile,omitempty"`
// Reachable is what can be reached on the machine now: every listening socket and every
// published container port. Only an adopted node reports it; it is what converging previews.
Reachable []Reach `json:"reachable,omitempty"`
@@ -273,16 +264,6 @@ type Held struct {
Changed string `json:"changed,omitempty"`
// Kept is where a file's original was kept.
Kept string `json:"kept,omitempty"`
// Facts is the found thing beside what the module declares — what a take compares (novox/hq
// ADR 0163): the host's own shape, carried as data and read by the preview.
Facts map[string]any `json:"facts,omitempty"`
}
// A Stray is a container a machine runs that the mesh neither wrote nor holds (ADR 0163).
type Stray struct {
Kind string `json:"kind"`
Name string `json:"name"`
Detail string `json:"detail,omitempty"`
}
// Reach is one thing reachable on the machine: a listening socket, or a published container port.