Compare commits

...
Author SHA1 Message Date
mesh-admin 6215ff0760 Merge pull request 'A machine says which networks it routes, and its filter forwards them' (#133) from feat/a-machine-says-which-networks-it-routes into main 2026-09-28 19:31:06 +00:00
jschoubben 54812306be A machine says which networks it routes, and its filter forwards them
The derived filter denies forwarding by default and then allows the container
runtime's two default pools, named in this code with a comment saying a machine
configured otherwise needs to say so -- and no way to say it. So the filter was
right on a machine using the defaults and silently wrong on any other.

Measured today: flipping a workstation to the derived filter cut egress for five
of its container networks and for every network its test beds create, because
those come from ranges the defaults do not cover. Nothing reported a fault; the
guests just could not reach anything, while the machine reported it had applied
what it was told.

A node-level fact beside the public domain, because the machine routes them and
the module that loads the filter may be replaced. Added to the defaults, never
replacing them. Their guests also keep address and name service, without which a
network does not work at all, and the converge preview now says what a machine
routes instead of leaving it to a sentence about what it cannot preview.
2026-09-28 21:29:10 +02:00
mesh-admin ce9e20fbbc Merge pull request 'A push raises what a module hears, not only a start' (#132) from fix/a-push-raises-what-a-module-hears into main 2026-09-28 14:46:54 +00:00
jschoubben 878690697e A push raises what a module hears, not only a start
A module's declaration and its consumer are derived from the same records, and only one of them
followed a push: the consumers were raised when the control plane started serving, so a module that
gained a `consumes` was sent a declaration it could act on and a consumer that never delivered the
event — with nothing anywhere saying the two disagreed. Found on review: the catalogue's own replay
subscription was recorded, granted and never delivered.

Everything the raise does is idempotent, so a push may do it.
2026-09-28 16:46:42 +02:00
mesh-admin ad97297576 Merge pull request 'The seat declares the facts its holder states' (#131) from fix/the-seat-declares-the-facts-its-holder-states into main 2026-09-28 14:37:05 +00:00
jschoubben 683b1ed693 The seat declares the facts its holder states
The grant permitted the control plane to state what it applied and the seat said nothing about it, so
the check that every derived subscription has an owner found the catalogue subscribing to a subject
nothing publishes — which is exactly the fault that check exists for, pointed at me.

A seat carries the protocol of its role (novox/hq ADR 0129), so the facts are the mesh-controller
seat's `emits`. That is also what lets another module declare it consumes them. No accepts, so no work
queue is raised for the seat — only what its holder may say. The agreement test now holds all three
places to one another: the seat, the grant, and the words the mesh states them with.
2026-09-28 16:36:45 +02:00
mesh-admin 04f9f378b0 Merge pull request 'Two faults found on review' (#130) from fix/review-two-small-faults into main 2026-09-28 14:32:43 +00:00
jschoubben 1c3f44a526 Two faults found on review
A second Accept value would have overwritten the first, because the header was Set per value rather
than Added. One value is all any caller passes today, so nothing was wrong — but a helper that
quietly keeps only the last of what it was given is a trap for whoever passes two.

And a replayed announcement that could not be written was published as an empty body: a fact on the
mesh that says nothing, which the reader can only log and drop. It is now said and skipped, because a
body that cannot be marshalled is this program's fault rather than the bus's.
2026-09-28 16:32:41 +02:00
mesh-admin 89e152dfe2 Merge pull request 'The mesh says what it applied, and the replay has an address it may use' (#129) from feat/the-mesh-says-what-it-applied into main 2026-09-28 14:07:20 +00:00
jschoubben 1ebad3786c The mesh says what it applied, and the replay has an address it may use
The pipeline was observable from a merge to an artifact and went dark where it touched a machine: a
node's report is control traffic only the control plane reads, so nothing said which version a
machine runs, or that it refused to (novox/hq ADR 0134). The control plane now states both under the
seat it holds — a role's events belong to the role and keep their address when the holder is
replaced — and only when the report is news, because a machine reconciles every minute and a fact per
report would be a fact per minute per machine.

Whether a report is news is the store's answer: it holds the previous one, so the listener returns it
and the server states the fact. That also gives the catch-up replay a subject the controller may
publish: it was published as a module's event from a module called "control-plane", which does not
exist, so the controller's own account refused it and every catalogue that asked what it missed was
answered with nothing.
2026-09-28 16:07:18 +02:00
mesh-admin f2f526a60a Merge pull request 'A module's name may contain a dot, so the derived step adds none' (#128) from fix/a-modules-name-may-contain-a-dot into main 2026-09-28 13:44:21 +00:00
jschoubben 4b4c7e0e0d A module's name may contain a dot, so the derived step adds none
A resource's id is `<module>.<its own id>` and a module's name may itself contain a dot — novox.be is
one — so the owner of a resource is everything before the *last* dot. The preparation step's id used a
dot, which made its owner unreadable by that rule; it uses a hyphen, and the id says what it belongs
to whichever way a reader splits it.
2026-09-28 15:44:18 +02:00
mesh-admin cec792ce9d Merge pull request 'A manifest HEAD says what it accepts, or the registry answers 404' (#127) from fix/a-manifest-head-says-what-it-accepts into main 2026-09-28 11:02:10 +00:00
jschoubben 338d033632 A manifest HEAD says what it accepts, or the registry answers 404
The check that skips copying a base the mesh already holds asked with no Accept header, and a
registry answers a manifest only in a media type the caller named: the same digest answered 200 with
the manifest types and 404 without them. So the builder concluded it held nothing, copied every
vendor base again, and exhausted the public hub's pull limit a second time today.

The test could not have caught it, because the fake registry answered a manifest HEAD regardless of
Accept — more permissive than the thing it stands in for. It is now as strict as a real registry, and
fails without the fix.
2026-09-28 13:02:08 +02:00
mesh-admin 1be926cec4 Merge pull request 'The control plane prepares its own state, like any module' (#126) from feat/the-control-plane-prepares-its-own-state into main 2026-09-28 10:51:30 +00:00
jschoubben 2134768dfe The control plane prepares its own state, like any module
Now that every parser on the mesh knows the word, the control plane's manifest says it. Its schema
stops being a special case: the mesh derives the step from its own resource and gates its server on
it, which is the failure of novox/hq 04-ISSUES/133 closed by the mechanism rather than by a
hand-written step in one manifest.
2026-09-28 12:51:27 +02:00
mesh-admin ef825688ee Merge pull request 'The control plane learns 'prepares' one release before its manifest uses it' (#125) from fix/the-word-ships-before-the-manifest-uses-it into main 2026-09-28 10:49:36 +00:00
jschoubben 77a14360df The control plane learns 'prepares' one release before its manifest uses it
A manifest word has to reach every parser before a manifest carries it. The builder refused
mesh-controller's manifest with `unknown field "prepares"` until it was rebuilt; then the running
control plane could not read the manifest in the build result either, so the build was recorded with
no module and the version never moved. Strict parsing is deliberate (novox/hq 04-ISSUES/003), so the
word ships first and a manifest uses it next: this takes `prepares` back out of the control plane's
own manifest, leaving the code that understands it, and the manifest says it again once this is
running everywhere.
2026-09-28 12:49:33 +02:00
mesh-admin 9be2fb4750 Merge pull request 'A version prepares its state before it runs' (#124) from feat/a-version-prepares-its-state into main 2026-09-28 10:43:17 +00:00
jschoubben 5a963aec10 A version prepares its state before it runs
The mesh derives the preparation from the module's own resource instead of each module hand-writing
a step beside it (novox/hq ADR 0135). A manifest says one word — `prepares` — and the mesh runs that
module's own program in its preparation mode, in the module's own context: the same image, the same
environment, the same mounts, because it is the same code. A published port and a fixed address are
taken away rather than copied, since the version being replaced still holds them.

One word for every kind of module: a Go binary receives `prepare` as its argument, a bundle receives
it through the runtime whose entry takes the same word. The control plane answers it like anything
else — its own schema stops being a special case, and its hand-written step is gone.
2026-09-28 12:43:15 +02:00
mesh-admin 45d1c28a28 Merge pull request 'The control plane migrates before it serves' (#123) from fix/the-control-plane-migrates-before-it-serves into main 2026-09-28 08:27:44 +00:00
mesh-admin da394b45e6 Merge pull request 'Work slower than the window says so, and one address is the bus's' (#122) from fix/work-longer-than-the-window-says-so into main 2026-09-28 07:51:52 +00:00
jschoubben 2f3bfda8c0 Work slower than the window says so, and one address is the bus's
Three faults the mesh's own logs showed this morning. A handler that outlives the acknowledgement
window was handed its message again while it was still working: acting on a merge builds modules,
minutes against a thirty-second window, so one merge ran the whole catalogue five times over. The
transport now says the work is in progress while it runs, which is where the window belongs.

Everything the mesh hands out — a token, a membership, a person's credential — took its address
from the enrolment setting, which on a mesh that has moved still names the broker it moved from:
the first person issued after the move was handed the retired broker's port. There is one bus, and
its address is the one the control plane is connected to.

And `operator issue` documented an argument order its parser refused.
2026-09-28 09:51:50 +02:00
39 changed files with 1149 additions and 123 deletions
+1 -1
View File
@@ -88,7 +88,7 @@ func reportsReaching(t *testing.T, open *stores, reachable []link.Reach, held ..
if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil { if err := open.inventory.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{ if _, err := (link.Enrolment{Inventory: open.inventory}).Heard(ctx, link.Report{
Node: "anchor", Applied: []string{"hello-web.x"}, Declared: digestOf(body), Node: "anchor", Applied: []string{"hello-web.x"}, Declared: digestOf(body),
Firewall: "ufw", Held: held, Reachable: reachable, Firewall: "ufw", Held: held, Reachable: reachable,
}); err != nil { }); err != nil {
+41 -1
View File
@@ -7,6 +7,7 @@ import (
"errors" "errors"
"flag" "flag"
"fmt" "fmt"
"net"
"slices" "slices"
"sort" "sort"
"strings" "strings"
@@ -309,7 +310,7 @@ func converge(ctx context.Context, open *stores, node string, yes bool, digest s
return "", err return "", err
} }
derived := derivedFilter{rules: rules, foundation: with.Foundation, mesh: with.Mesh, derived := derivedFilter{rules: rules, foundation: with.Foundation, mesh: with.Mesh,
outward: plan.PublicDomain != ""} outward: plan.PublicDomain != "", routed: with.Routed}
preview, saw := previewOf(node, reported, derived, plan, taken, filter, runs[filter]) preview, saw := previewOf(node, reported, derived, plan, taken, filter, runs[filter])
preview += "\n\n preview " + saw preview += "\n\n preview " + saw
if !yes { if !yes {
@@ -414,6 +415,17 @@ func previewOf(node string, reported inventory.Adoption, derived derivedFilter,
b.WriteString(" not previewed: traffic the machine routes that is not a published port " + b.WriteString(" not previewed: traffic the machine routes that is not a published port " +
"(a tunnel, NAT in the found firewall) — the derived filter drops it unless a module " + "(a tunnel, NAT in the found firewall) — the derived filter drops it unless a module " +
"declares it\n") "declares it\n")
// What it routes, said rather than left to the sentence above (novox/hq ADR 0137). The filter
// forwards the container runtime's own default pools without being told; anything else is this
// list, and a machine whose guests live outside those pools and names none of them loses their
// egress at the flip, silently, which is how this was found.
if len(derived.routed) > 0 {
b.WriteString(fmt.Sprintf(" it routes: %s — kept forwarded, and their guests keep "+
"address and name service\n", strings.Join(derived.routed, ", ")))
} else {
b.WriteString(" it routes: nothing said, so only the container runtime's own default " +
"pools are forwarded — `node networks <node> <cidr>...` if its guests live elsewhere\n")
}
isTaken := map[string]bool{} isTaken := map[string]bool{}
for _, m := range taken { for _, m := range taken {
@@ -473,6 +485,9 @@ type derivedFilter struct {
// mesh is every address on the private network; outward says the machine faces outside. // mesh is every address on the private network; outward says the machine faces outside.
mesh []string mesh []string
outward bool outward bool
// routed is the networks this machine says it routes for what it hosts (novox/hq ADR 0137):
// their guests keep address and name service, and what they send onward keeps being forwarded.
routed []string
} }
// closesOutside is what a narrowing from everywhere to the private network is called: it closes. // closesOutside is what a narrowing from everywhere to the private network is called: it closes.
@@ -498,6 +513,13 @@ func (d derivedFilter) fate(r inventory.Reach) string {
return "stays open — the mesh's own, from anywhere" return "stays open — the mesh's own, from anywhere"
} }
} }
// A guest on a network this machine routes asks it for an address and for names, and those two
// arrive here (novox/hq ADR 0137). Matched on the listener's own address: a resolver bound to a
// bridge in one of those networks is the one its guests ask.
if within(r.Address, d.routed) &&
((r.Protocol == "udp" && (r.Port == 53 || r.Port == 67)) || (r.Protocol == "tcp" && r.Port == 53)) {
return "stays open — address and name service for a network this machine routes"
}
for _, rule := range d.rules { for _, rule := range d.rules {
if rule.Port != r.Port || rule.Protocol != r.Protocol { if rule.Port != r.Port || rule.Protocol != r.Protocol {
continue continue
@@ -520,6 +542,24 @@ func (d derivedFilter) fate(r inventory.Reach) string {
return "WILL CLOSE — no module assigned here declares it" return "WILL CLOSE — no module assigned here declares it"
} }
// within says whether an address the machine reported sits inside one of the networks it routes.
func within(address string, networks []string) bool {
ip := net.ParseIP(strings.Trim(address, "[]"))
if ip == nil {
return false
}
for _, n := range networks {
_, block, err := net.ParseCIDR(n)
if err != nil {
continue
}
if block.Contains(ip) {
return true
}
}
return false
}
// countReachable is how much of a node's account of itself names something off the machine. // countReachable is how much of a node's account of itself names something off the machine.
// Loopback is left out for the same reason the preview leaves it out: nothing outside reaches it, // Loopback is left out for the same reason the preview leaves it out: nothing outside reaches it,
// so a report of loopback alone says nothing about what the filter would close. // so a report of loopback alone says nothing about what the filter would close.
+8 -1
View File
@@ -76,7 +76,10 @@ func run() error {
return pinCommand(ctx, args[1:], true) return pinCommand(ctx, args[1:], true)
case "unpin": case "unpin":
return pinCommand(ctx, args[1:], false) return pinCommand(ctx, args[1:], false)
case "migrate": // `prepare` is how the mesh asks any module to bring its state to the shape this version needs
// (novox/hq ADR 0135), and the control plane answers it the same way as everything else — its
// own schema is not a special case. `migrate` remains the word a person types.
case "prepare", "migrate":
return migrate(ctx) return migrate(ctx)
case "node": case "node":
return nodeCommand(ctx, args[1:]) return nodeCommand(ctx, args[1:])
@@ -138,12 +141,16 @@ func usage() {
fmt.Fprint(os.Stderr, `mesh-controller — the control plane fmt.Fprint(os.Stderr, `mesh-controller — the control plane
migrate bring each context's schema up to date migrate bring each context's schema up to date
prepare the same, asked the way the mesh asks any module (ADR 0135)
node add <name> [--adopted] create a node record; --adopted: the machine is in use node add <name> [--adopted] create a node record; --adopted: the machine is in use
node list the nodes this mesh knows about node list the nodes this mesh knows about
node show <name> what one machine reported it can do, and why node show <name> what one machine reported it can do, and why
node public-domain <name> the domain it composes its routed names under node public-domain <name> the domain it composes its routed names under
node public-domain <name> <d> ...set it to d node public-domain <name> <d> ...set it to d
node public-domain <name> --clear ...it faces the outside no longer node public-domain <name> --clear ...it faces the outside no longer
node networks <name> the networks it routes for what it hosts
node networks <name> <cidr>... ...set them; its filter forwards these too
node networks <name> --clear ...only the container runtime's own
token issue --node <name> a one-time right to join, for an existing record token issue --node <name> a one-time right to join, for an existing record
token issue --new <name> create the record and issue for it token issue --new <name> create the record and issue for it
token issue ... --adopted ...for a machine in use, which joins adopted token issue ... --adopted ...for a machine in use, which joins adopted
+69 -1
View File
@@ -21,7 +21,8 @@ import (
func nodeCommand(ctx context.Context, args []string) error { func nodeCommand(ctx context.Context, args []string) error {
if len(args) == 0 { if len(args) == 0 {
return errors.New("node add <name>, node list, node show <name>, or " + publicDomainUsage) return errors.New("node add <name>, node list, node show <name>, " + publicDomainUsage +
", or " + networksUsage)
} }
open, err := openStores(ctx) open, err := openStores(ctx)
if err != nil { if err != nil {
@@ -66,6 +67,12 @@ func nodeCommand(ctx context.Context, args []string) error {
// because the damage is already done by the time it prints. // because the damage is already done by the time it prints.
return publicDomain(ctx, inv, args[1:]) return publicDomain(ctx, inv, args[1:])
case "networks":
// The networks this machine routes for what it hosts (novox/hq ADR 0137): what the derived
// filter must keep forwarding, beyond the container runtime's own default pools which it
// allows without being told. Reports with no argument, for the same reason the domain does.
return nodeNetworks(ctx, inv, args[1:])
case "account": case "account":
// The operator's login on this machine (novox/hq to-be 29): what a home-scoped file is // The operator's login on this machine (novox/hq to-be 29): what a home-scoped file is
// owned by and which account `ssh <node>` uses. Reports with no argument; sets with one; // owned by and which account `ssh <node>` uses. Reports with no argument; sets with one;
@@ -144,6 +151,67 @@ func nodeAccount(ctx context.Context, inv *inventory.Inventory, positionals []st
return nil return nil
} }
const networksUsage = "node networks <name> — what it routes now; " +
"<name> <cidr>... to set them; <name> --clear to route only the container runtime's own"
// nodeNetworks reads, sets or clears the networks a machine routes for what it hosts.
//
// The same three forms as the domain above, and the read-shaped one reports rather than clearing,
// for the same reason: this list is what keeps a machine's guests reaching anything, and losing it
// by asking a question is not a mistake anybody can see afterwards.
func nodeNetworks(ctx context.Context, inv *inventory.Inventory, args []string) error {
set := flag.NewFlagSet("node networks", flag.ContinueOnError)
clear := set.Bool("clear", false,
"route only the container runtime's own default pools, as a machine that has said nothing does")
positionals, err := parseAround(set, args)
if err != nil {
return err
}
if len(positionals) == 0 {
return errors.New(networksUsage)
}
node := positionals[0]
switch {
case *clear && len(positionals) > 1:
return fmt.Errorf("give %s networks or --clear, not both: %q and --clear say opposite "+
"things and the mesh will not choose between them", node, strings.Join(positionals[1:], " "))
case *clear:
if err := inv.SetRoutedNetworks(ctx, node, nil); err != nil {
return err
}
fmt.Printf("%s routes only the container runtime's own default pools\n", node)
fmt.Printf(" run `push %s` to send its filter\n", node)
return nil
case len(positionals) > 1:
if err := inv.SetRoutedNetworks(ctx, node, positionals[1:]); err != nil {
return err
}
fmt.Printf("%s routes %s\n", node, strings.Join(positionals[1:], ", "))
fmt.Printf(" its filter forwards them, and their guests keep address and name service\n")
fmt.Printf(" run `push %s` to send it\n", node)
return nil
default:
if _, err := inv.NodeByName(ctx, node); err != nil {
return err
}
networks, err := inv.RoutedNetworksOf(ctx, node)
if err != nil {
return err
}
if len(networks) == 0 {
fmt.Printf("%s routes only the container runtime's own default pools\n", node)
fmt.Printf(" `node networks %s <cidr>...` if its guests live elsewhere\n", node)
return nil
}
fmt.Printf("%s routes %s\n", node, strings.Join(networks, ", "))
return nil
}
}
const publicDomainUsage = "node public-domain <name> — what it is now; " + const publicDomainUsage = "node public-domain <name> — what it is now; " +
"<name> <domain> to set it; <name> --clear to take it away" "<name> <domain> to set it; <name> --clear to take it away"
+7 -3
View File
@@ -189,13 +189,17 @@ func readPrivateKey(path string) (string, error) {
func personIssue(ctx context.Context, args []string) error { func personIssue(ctx context.Context, args []string) error {
set := flag.NewFlagSet("operator issue", flag.ContinueOnError) set := flag.NewFlagSet("operator issue", flag.ContinueOnError)
invokes := set.String("invokes", "", "the tools this person may call, comma-separated, or * for every one") invokes := set.String("invokes", "", "the tools this person may call, comma-separated, or * for every one")
if err := set.Parse(args); err != nil { // Flags on either side of the name, because the usage this command prints puts them after it —
// and the standard parser stops at the first thing that is not a flag, so the order the command
// documents was the one order it refused (2026-09-28).
positionals, err := parseAround(set, args)
if err != nil {
return err return err
} }
if set.NArg() != 1 { if len(positionals) != 1 {
return errors.New("operator issue <name> --invokes <tool,tool|*>") return errors.New("operator issue <name> --invokes <tool,tool|*>")
} }
name := set.Arg(0) name := positionals[0]
if *invokes == "" { if *invokes == "" {
return errors.New( return errors.New(
"say what this person may call: --invokes mesh-catalog.catalog_tools,gitea.repo_create, " + "say what this person may call: --invokes mesh-catalog.catalog_tools,gitea.repo_create, " +
+7 -1
View File
@@ -640,13 +640,19 @@ func renderingFor(ctx context.Context, open *stores, node string,
if err != nil { if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err return catalogue.Rendering{}, inventory.Node{}, err
} }
// The networks this machine routes for what it hosts, which the derived filter must forward for
// (novox/hq ADR 0137).
routed, err := inv.RoutedNetworksOf(ctx, node)
if err != nil {
return catalogue.Rendering{}, inventory.Node{}, err
}
return catalogue.Rendering{ return catalogue.Rendering{
BusMembership: memberships[node], BusMembership: memberships[node],
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports, Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
Certificate: certificate, Authority: authority, Mesh: private, Names: names, Certificate: certificate, Authority: authority, Mesh: private, Names: names,
Machines: machines, Machines: machines,
Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation, Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation,
Kept: kept, Adopted: record.Adopted, Kept: kept, Adopted: record.Adopted, Routed: routed,
Given: given, Taken: taken, Seats: seats, ArtifactStore: artifactStore, Built: built, Given: given, Taken: taken, Seats: seats, ArtifactStore: artifactStore, Built: built,
BusUsers: busUsers, BusUsers: busUsers,
}, record, nil }, record, nil
+7 -1
View File
@@ -162,7 +162,13 @@ func declare(ctx context.Context, args []string) error {
return err return err
} }
server, err := connectLink(ctx, nil, nil, nil) // **With the inventory, so the bus is raised** (novox/hq ADR 0134, design 30). A module's
// declaration and how it hears what it consumes move together: its consumer is derived from the
// same records this declaration is composed from. Raised only when the control plane started
// serving, a module that gained a `consumes` was sent a declaration it could act on and a
// consumer that never delivered the event — and nothing anywhere said the two disagreed
// (found on review, 2026-09-28). Everything the raise does is idempotent.
server, err := connectLink(ctx, inv, nil, nil)
if err != nil { if err != nil {
return err return err
} }
+9
View File
@@ -245,10 +245,12 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
// commit it already records — so it is rebuilt and its record left alone. Writing this commit as // commit it already records — so it is rebuilt and its record left alone. Writing this commit as
// its source would make it permanently behind a repository its manifest does not come from. // its source would make it permanently behind a repository its manifest does not come from.
var from, packaging []inventory.Entry var from, packaging []inventory.Entry
already := 0
for _, e := range entries { for _, e := range entries {
switch { switch {
case sourceIs(e.Source, m): case sourceIs(e.Source, m):
if e.Source.BuiltFrom == m.Commit { if e.Source.BuiltFrom == m.Commit {
already++
continue continue
} }
// **A merge older than the last look at the source is history, not a move.** The forge // **A merge older than the last look at the source is history, not a move.** The forge
@@ -264,6 +266,13 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
} }
} }
if len(from) == 0 && len(packaging) == 0 { if len(from) == 0 && len(packaging) == 0 {
// "Already built from it" and "nothing reads it" are different facts, and reading the first
// as the second sends somebody looking for a broken trigger when the mesh is up to date.
if already > 0 {
fmt.Printf("%s/%s merged into %s (%.8s); %d module(s) the mesh holds are already built "+
"from it\n", m.Owner, m.Repo, m.Base, m.Commit, already)
return nil
}
fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds reads it\n", fmt.Printf("%s/%s merged into %s (%.8s); nothing the mesh holds reads it\n",
m.Owner, m.Repo, m.Base, m.Commit) m.Owner, m.Repo, m.Base, m.Commit)
return nil return nil
+13
View File
@@ -66,6 +66,19 @@ func FromEnvironment() (Broker, error) {
if err != nil { if err != nil {
return Broker{}, err return Broker{}, err
} }
// **One bus, one address** (novox/hq ADR 0131). Everything the mesh hands out — a token, a
// machine's membership, a person's credential — must name the bus the control plane itself is
// connected to; the setting above predates the move and, on a mesh that has moved, still names
// the broker it moved from. The first person issued after the move was handed the retired
// broker's port and could not connect to anything (2026-09-28).
//
// Read from the credential rather than from a second setting somebody keeps in step: the
// control plane cannot be wrong about where it is connected.
if bus, on, err := OnNATS(); err == nil && on {
if where := strings.TrimPrefix(BareAddress(bus), "nats://"); where != "" {
address = where
}
}
return Broker{Address: address, Fingerprint: fingerprint}, nil return Broker{Address: address, Fingerprint: fingerprint}, nil
} }
+8
View File
@@ -181,6 +181,14 @@ func PermissionsFor(p Principal) (Permissions, error) {
for _, seat := range meshSeatsTheControllerUses { for _, seat := range meshSeatsTheControllerUses {
pub = append(pub, "mesh.seat."+seat+".accept.>") pub = append(pub, "mesh.seat."+seat+".accept.>")
} }
// **And what the mesh says it did** (novox/hq ADR 0134). The control plane states its own
// facts under the seat it holds, because a role's events belong to the role and keep their
// address while the holder is replaced. Named one by one rather than as a whole namespace:
// least authority, and a fact nothing states is authority nobody uses.
for _, event := range ControllerStates {
pub = append(pub, seatEventSubject(ControllerSeat, event))
}
// Every module's tools: **the control plane is the way in** (novox/hq ADR 0095). A person // Every module's tools: **the control plane is the way in** (novox/hq ADR 0095). A person
// or an agent asks through it and every question passes one process where an audit // or an agent asks through it and every question passes one process where an audit
// belongs — so it, alone among principals, may call any tool by name. The first `ask` on // belongs — so it, alone among principals, may call any tool by name. The first `ask` on
+44
View File
@@ -0,0 +1,44 @@
package broker_test
import (
"slices"
"testing"
"github.com/novox/mesh-controller/internal/broker"
"github.com/novox/mesh-controller/internal/catalogue"
"github.com/novox/mesh-controller/internal/link"
)
// The facts the control plane states are named twice — in the grant that permits them and in the code
// that states them — because `link` imports `broker` and the dependency cannot go the other way. So a
// test keeps them agreeing: a subject the grant omits is refused at the moment the mesh has something
// to say, and one the grant adds that nothing states is authority nobody uses.
//
// An external test package, because it may import both while neither imports the other.
func TestTheFactsTheGrantPermitsAreTheFactsTheMeshStates(t *testing.T) {
if broker.ControllerSeat != link.MeshControllerSeat {
t.Fatalf("the grant is written for the %q seat and the mesh states its facts under %q",
broker.ControllerSeat, link.MeshControllerSeat)
}
for _, event := range []string{link.KeyApplied, link.KeyRefused, link.KeyBuiltBefore} {
if !slices.Contains(broker.ControllerStates, event) {
t.Errorf("the mesh states %q and its account may not publish it", event)
}
}
if len(broker.ControllerStates) != 3 {
t.Errorf("the grant permits %v, which is more than the mesh states", broker.ControllerStates)
}
// **And the seat says it.** A seat carries the protocol of its role (novox/hq ADR 0129), so the
// facts the control plane states are the seat's `emits` — which is what lets anything else declare
// that it consumes them, and what the subject-agreement check reads to know they have an owner.
var declared []string
for _, seat := range catalogue.SeatsWithAProtocol() {
if seat.Name == broker.ControllerSeat {
declared = seat.Emits
}
}
if !slices.Equal(declared, broker.ControllerStates) {
t.Errorf("the %s seat emits %v and the grant permits %v", broker.ControllerSeat,
declared, broker.ControllerStates)
}
}
+11
View File
@@ -160,6 +160,17 @@ func Overlaps() []string {
// ack subject is derived from (nats.go: `$JS.ACK.<stream>.controller.>`). // ack subject is derived from (nats.go: `$JS.ACK.<stream>.controller.>`).
const ControllerName = "controller" const ControllerName = "controller"
// ControllerSeat is the role the control plane holds, and ControllerStates are the facts it states
// under it (novox/hq ADR 0134).
//
// **Written here as well as in `link`, and a test keeps them agreeing.** `link` imports `broker`, so
// `broker` cannot import `link`; a grant naming a subject the controller never publishes is authority
// nobody uses, and a controller publishing one the grant omits is refused at the moment it has
// something to say.
const ControllerSeat = "mesh-controller"
var ControllerStates = []string{"applied", "refused", "built-before"}
// ControllerFollows are the events the controller reacts to: the catalogue saying a module's // ControllerFollows are the events the controller reacts to: the catalogue saying a module's
// current version moved, and a catalogue that has just started saying it may have missed builds. // current version moved, and a catalogue that has just started saying it may have missed builds.
// //
+1 -1
View File
@@ -24,7 +24,7 @@ accounts {
jetstream: enabled jetstream: enabled
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
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.>"] } 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"] } 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"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
+1 -1
View File
@@ -186,7 +186,7 @@ func (r Registry) MirrorImage(ctx context.Context, from, repository string) (str
// on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base // on the first merge that rebuilt a whole catalogue (2026-09-28), and every module whose base
// lives there failed on a copy it did not need. // lives there failed on a copy it did not need.
if strings.HasPrefix(where.reference, "sha256:") { if strings.HasPrefix(where.reference, "sha256:") {
held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference) held, err := r.has(ctx, "http://"+r.Address+"/v2/"+repository+"/manifests/"+where.reference, manifestAccept)
if err != nil { if err != nil {
return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err) return "", fmt.Errorf("asking %s whether it holds %s: %w", r.Address, from, err)
} }
+8
View File
@@ -104,6 +104,14 @@ func (m *theMeshsRegistry) handler() http.Handler {
defer m.mu.Unlock() defer m.mu.Unlock()
switch { switch {
case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"): case r.Method == http.MethodHead && strings.Contains(r.URL.Path, "/manifests/"):
// **As strictly as a real registry.** A manifest is answered only in a media type the
// caller named; a request with no Accept is answered as if nothing were there. The fake
// used to answer regardless, which is why it could not catch a check that asked without
// one — and the mesh copied every base again (2026-09-28).
if !strings.Contains(r.Header.Get("Accept"), "manifest") && !strings.Contains(r.Header.Get("Accept"), "index") {
w.WriteHeader(http.StatusNotFound)
return
}
if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok { if _, ok := m.manifests[r.URL.Path[strings.LastIndex(r.URL.Path, "/")+1:]]; ok {
w.WriteHeader(http.StatusOK) w.WriteHeader(http.StatusOK)
} else { } else {
+13 -1
View File
@@ -122,11 +122,23 @@ func (r Registry) PublishArchive(ctx context.Context, repository string, body []
return final, nil return final, nil
} }
func (r Registry) has(ctx context.Context, url string) (bool, error) { // has is whether this registry already holds what is at that URL.
//
// **A manifest HEAD must say what it accepts.** A registry answers a manifest request only in a media
// type the caller named, and a bare HEAD — no Accept at all — is answered 404 for a manifest it holds
// perfectly well. Measured against the mesh's own registry (2026-09-28): the same digest answered 200
// with the manifest media types and 404 without them, so a check written without them concluded the
// registry held nothing, copied every base again, and exhausted the public hub's pull limit. A blob
// needs no Accept, which is why this went unnoticed: the same helper was right for blobs and wrong
// for manifests.
func (r Registry) has(ctx context.Context, url string, accept ...string) (bool, error) {
request, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil) request, err := http.NewRequestWithContext(ctx, http.MethodHead, url, nil)
if err != nil { if err != nil {
return false, err return false, err
} }
for _, media := range accept {
request.Header.Add("Accept", media)
}
response, err := r.client().Do(request) response, err := r.client().Do(request)
if err != nil { if err != nil {
return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err) return false, fmt.Errorf("cannot reach the registry at %s: %w", r.Address, err)
+94 -1
View File
@@ -124,6 +124,11 @@ type Rendering struct {
// nothing on this node keeps them, or the mesh has no operator key. // nothing on this node keeps them, or the mesh has no operator key.
Kept *KeptExport Kept *KeptExport
// Routed is the networks this machine routes for what it hosts, beyond the container runtime's
// own default pools, which the filter allows without being told (novox/hq ADR 0137). A node-level
// fact: the machine routes them, and the module that loads the filter may be replaced.
Routed []string
// Foundation is the ports the mesh itself needs reachable on every machine, which no module // Foundation is the ports the mesh itself needs reachable on every machine, which no module
// declares because the foundation is not a module (novox/hq 04-ISSUES/051 and 052). The broker // declares because the foundation is not a module (novox/hq 04-ISSUES/051 and 052). The broker
// is the one that matters: a machine dials it to enrol, and a firewall derived only from // is the one that matters: a machine dials it to enrol, and a firewall derived only from
@@ -345,7 +350,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
if err != nil { if err != nil {
return nil, err return nil, err
} }
filtering := AsNftables(rules, with.Mesh, r.PublicDomain != "", with.Foundation) filtering := AsNftables(rules, with.Mesh, r.PublicDomain != "", with.Foundation, with.Routed)
var out []map[string]any var out []map[string]any
for _, m := range r.Modules { for _, m := range r.Modules {
@@ -623,6 +628,9 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
// one of them as its environment without saying so (ADR 0086, issue 041). // one of them as its environment without saying so (ADR 0086, issue 041).
secretFiles := secretFilesOf(resources) secretFiles := secretFilesOf(resources)
// Which of this module's resources its preparation runs before, if it prepares anything.
prepareBefore := preparationTarget(m)
for _, unsettled := range resources { for _, unsettled := range resources {
resource, err := ApplySettings(unsettled, with.Settings[m.Module]) resource, err := ApplySettings(unsettled, with.Settings[m.Module])
if err != nil { if err != nil {
@@ -708,6 +716,18 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil { if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
copied["reload-on"] = renamed copied["reload-on"] = renamed
} }
// **A version prepares its state before it runs** (novox/hq ADR 0135). Derived from the
// module's own resource rather than declared beside it: what prepares the state is the
// module's own code, so what it is given has to be what that code is given — and a
// second resource written by hand is a second copy to drift from the first. Placed
// immediately before it, because a run-once step stops everything the declaration
// places after it (ADR 0052), which is how a version whose preparation failed does not
// serve.
if prepareBefore != "" && fmt.Sprint(resource["id"]) == prepareBefore {
step := prepared(copied)
owner[fmt.Sprint(step["id"])] = m.Module
out = append(out, step)
}
owner[fmt.Sprint(copied["id"])] = m.Module owner[fmt.Sprint(copied["id"])] = m.Module
out = append(out, copied) out = append(out, copied)
} }
@@ -1610,3 +1630,76 @@ func atMachinePort(serves map[string]any, module string, ports map[string]map[in
func AtPublishedPort(values map[string]any, module string, published map[int]int) map[string]any { func AtPublishedPort(values map[string]any, module string, published map[int]int) map[string]any {
return atMachinePort(values, module, map[string]map[int]int{module: published}) return atMachinePort(values, module, map[string]map[int]int{module: published})
} }
// PreparationArgument is how the mesh asks a module to prepare its state: one word, to the module's
// own program, whatever that program is (novox/hq ADR 0135).
//
// **One word for every kind of module.** A module built as a Go binary receives it as its argument;
// one built as a bundle receives it through the runtime, whose entry takes the same word. So the
// mesh has one way of asking and a module has one way of answering, and neither learns the other's
// shape.
const PreparationArgument = "prepare"
// preparationTarget is the resource a module's preparation runs before: its own workload.
//
// The first container carrying an artifact this module built, and not itself a step — that is the
// thing that runs the module's code, and therefore the thing whose state must be ready. Empty when
// the module prepares nothing, or when nothing it declares could run its code.
//
// **A module with two own workloads gates the first of them.** Five modules in the catalogue declare
// more than one container of their own, none of them preparing anything today. If one ever does and
// its second workload shares the state, the gate is in front of the first — stated here because the
// alternative is a field asking an author to restate what the mesh can see.
func preparationTarget(m Manifest) string {
if !m.Prepares {
return ""
}
for _, r := range m.Resources {
if fmt.Sprint(r["type"]) != "container" || !ownArtifact(r, m.Module) {
continue
}
if once, _ := r["run-once"].(bool); once {
continue
}
return fmt.Sprint(r["id"])
}
return ""
}
// ownArtifact is whether a resource runs something this module built, in either spelling a manifest
// may be in: naming the artifact, before a build resolved it, or carrying the reference a build
// recorded — this mesh's own store, under this module's name.
func ownArtifact(resource map[string]any, module string) bool {
if named, _ := resource["artifact"].(string); named != "" {
return true
}
image, _ := resource["image"].(string)
return strings.HasPrefix(image, ArtifactStoreScheme+module+"/")
}
// prepared is the module's own resource as the step that prepares its state: the same image, the same
// context, run to completion with the mesh's preparation argument.
//
// Three things are taken away rather than copied, each because the step runs while the version it
// prepares for is still running. A published port cannot be bound twice, and a step that tried would
// fail for a reason that has nothing to do with the state. A fixed address cannot be held twice, for
// the same reason. And a cadence is what a step is the opposite of: a container runs once and gates,
// or on a schedule, or stays up, never two (ADR 0053).
func prepared(from map[string]any) map[string]any {
step := map[string]any{}
for k, v := range from {
step[k] = v
}
// **A hyphen, not a dot.** A resource's id is `<module>.<its own id>`, and a module's name may
// itself contain a dot (`novox.be`), so the module is everything before the *last* dot — which
// only works if what the mesh derives adds no dot of its own.
step["id"] = fmt.Sprint(from["id"]) + "-prepare"
step["name"] = fmt.Sprint(from["name"]) + "-prepare"
step["run-once"] = true
step["args"] = []any{PreparationArgument}
delete(step, "ports")
delete(step, "ip")
delete(step, "schedule")
delete(step, "reload-on")
return step
}
+33 -1
View File
@@ -230,7 +230,7 @@ const SSHPort = 22
// It is a floor for the same reason ssh is. A machine nobody can reach is a machine nobody can // It is a floor for the same reason ssh is. A machine nobody can reach is a machine nobody can
// repair; a machine the mesh cannot reach is a machine the mesh cannot manage. Neither is a thing // repair; a machine the mesh cannot reach is a machine the mesh cannot manage. Neither is a thing
// any module asks for, and neither may be derived away. // any module asks for, and neither may be derived away.
func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int) string { func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int, routed []string) string {
var b strings.Builder var b strings.Builder
b.WriteString("# Computed by the mesh from what is assigned to this node.\n") b.WriteString("# Computed by the mesh from what is assigned to this node.\n")
b.WriteString("# Edits are lost on the next declaration; change a module's listens instead.\n\n") b.WriteString("# Edits are lost on the next declaration; change a module's listens instead.\n\n")
@@ -252,6 +252,21 @@ func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int) str
b.WriteString("\t\ticmp type echo-request accept\n") b.WriteString("\t\ticmp type echo-request accept\n")
b.WriteString("\t\ticmpv6 type { echo-request, nd-neighbor-solicit, nd-neighbor-advert, nd-router-advert } accept\n") b.WriteString("\t\ticmpv6 type { echo-request, nd-neighbor-solicit, nd-neighbor-advert, nd-router-advert } accept\n")
// **What a machine it routes for must be able to ask it** (novox/hq ADR 0137). A guest on one of
// these networks gets its address and its names from this machine, over the bridge it is on, and
// those two questions arrive at the input chain like any other. Denied, the guest never gets an
// address and never resolves a name — which is not "a closed port" but a network that does not
// work at all, and it is this machine's own guest asking.
//
// Only these ports, and only for a network that was named: everything else a guest might want
// from its host is a port somebody declares, like every other port on this machine.
for _, network := range routed {
family := saddrFamily(network)
b.WriteString("\t\t# address and name service for a network this machine routes\n")
b.WriteString(fmt.Sprintf("\t\t%s saddr %s udp dport { 53, 67 } accept\n", family, network))
b.WriteString(fmt.Sprintf("\t\t%s saddr %s tcp dport 53 accept\n", family, network))
}
// **ssh, always, and not because a module asked.** // **ssh, always, and not because a module asked.**
// //
// Every other line in this chain is derived from what is assigned here, which is the whole // Every other line in this chain is derived from what is assigned here, which is the whole
@@ -375,6 +390,13 @@ func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int) str
b.WriteString(fmt.Sprintf("\t\t# %s\n", network.why)) b.WriteString(fmt.Sprintf("\t\t# %s\n", network.why))
b.WriteString(fmt.Sprintf("\t\tip saddr %s accept\n", network.cidr)) b.WriteString(fmt.Sprintf("\t\tip saddr %s accept\n", network.cidr))
} }
// And what this machine says it routes beyond them (novox/hq ADR 0137). Added to the defaults
// above, never replacing them: a machine that names one range has not stopped hosting whatever
// was already on the runtime's own.
for _, network := range routed {
b.WriteString("\t\t# a network this machine routes for what it hosts\n")
b.WriteString(fmt.Sprintf("\t\t%s saddr %s accept\n", saddrFamily(network), network))
}
if len(rules) > 0 { if len(rules) > 0 {
b.WriteString("\n") b.WriteString("\n")
@@ -442,6 +464,16 @@ var runtimeNetworks = []struct{ cidr, why string }{
{"192.168.128.0/17", "the networks its compose files are given"}, {"192.168.128.0/17", "the networks its compose files are given"},
} }
// saddrFamily is the match a network's family is written with: `ip saddr` or `ip6 saddr`. One match
// for both families is a syntax error, and a ruleset that does not load is a machine filtering
// nothing while its service reports a fault — the same reason byFamily below exists.
func saddrFamily(network string) string {
if strings.Contains(network, ":") {
return "ip6"
}
return "ip"
}
// byFamily splits addresses into the two nftables understands separately. // byFamily splits addresses into the two nftables understands separately.
// //
// `ip saddr` and `ip6 saddr` are different matches, and one set holding both families is a syntax // `ip saddr` and `ip6 saddr` are different matches, and one set holding both families is a syntax
@@ -19,7 +19,7 @@ func TestTheBrokersPortIsOpenedThoughNoModuleDeclaresIt(t *testing.T) {
// A machine on the private network, with one ordinary module rule, and nothing that mentions // A machine on the private network, with one ordinary module rule, and nothing that mentions
// the broker — which is every machine. // the broker — which is every machine.
rules := []Rule{{Port: 8080, From: FromMesh, Because: []string{"some-module"}}} rules := []Rule{{Port: 8080, From: FromMesh, Because: []string{"some-module"}}}
out := AsNftables(rules, []string{"10.42.0.1"}, false, []int{brokerPort}) out := AsNftables(rules, []string{"10.42.0.1"}, false, []int{brokerPort}, nil)
if !strings.Contains(out, "tcp dport 5671 accept") { if !strings.Contains(out, "tcp dport 5671 accept") {
t.Fatalf("the broker's port is not opened, so no machine could enrol:\n%s", out) t.Fatalf("the broker's port is not opened, so no machine could enrol:\n%s", out)
@@ -48,7 +48,7 @@ func TestTheBrokersPortIsOpenedThoughNoModuleDeclaresIt(t *testing.T) {
// And a mesh that was never told about a broker still gets a ruleset, rather than an empty one or // And a mesh that was never told about a broker still gets a ruleset, rather than an empty one or
// a panic. A control plane in that state cannot issue tokens either, which is where it surfaces. // a panic. A control plane in that state cannot issue tokens either, which is where it surfaces.
func TestNoBrokerMeansNoFoundationRuleRatherThanNoRuleset(t *testing.T) { func TestNoBrokerMeansNoFoundationRuleRatherThanNoRuleset(t *testing.T) {
out := AsNftables(nil, []string{"10.42.0.1"}, false, nil) out := AsNftables(nil, []string{"10.42.0.1"}, false, nil, nil)
if !strings.Contains(out, "table inet mesh") { if !strings.Contains(out, "table inet mesh") {
t.Fatalf("no ruleset at all:\n%s", out) t.Fatalf("no ruleset at all:\n%s", out)
} }
@@ -0,0 +1,99 @@
package catalogue
import (
"strings"
"testing"
)
// A machine's own guests keep working when the filter it is given denies by default.
//
// **The forward chain allowed the container runtime's two default pools and nothing else** — named
// in this package's code with a comment saying a machine configured otherwise "needs this to say
// so", and no way to say it (novox/hq ADR 0137). Measured on a workstation on 2026-09-28: the flip
// to the derived filter cut egress for five of its container networks, allocated from ranges those
// two defaults do not cover, and for every network its test beds create. Nothing reported a fault.
// The guests simply could not reach anything, and the machine went on saying it had applied what it
// was told.
func TestWhatAMachineSaysItRoutesKeepsBeingForwarded(t *testing.T) {
const beds = "10.0.0.0/8"
ruleset := AsNftables(nil, []string{"10.10.0.1"}, false, nil, []string{beds})
forward := chainOf(t, ruleset, "forward")
if !strings.Contains(forward, "ip saddr "+beds+" accept") {
t.Errorf("the forward chain does not accept what the machine says it routes (%s):\n%s",
beds, forward)
}
// The defaults stay. A machine that names one range has not stopped hosting whatever was
// already on the runtime's own pools, and losing those would trade one silent breakage for
// another.
for _, network := range runtimeNetworks {
if !strings.Contains(forward, "ip saddr "+network.cidr+" accept") {
t.Errorf("naming a network dropped the runtime's own %s:\n%s", network.cidr, forward)
}
}
// And its guests can still ask this machine the two questions that make a network usable at
// all: what is my address, and what is that name.
input := chainOf(t, ruleset, "input")
for _, want := range []string{
"ip saddr " + beds + " udp dport { 53, 67 } accept",
"ip saddr " + beds + " tcp dport 53 accept",
} {
if !strings.Contains(input, want) {
t.Errorf("the input chain is missing %q, so a guest on %s gets no address and "+
"resolves no name:\n%s", want, beds, input)
}
}
}
// A machine that says nothing is filtered exactly as it was before this existed.
//
// The change has to be additive on every machine already converged: novox has been carrying this
// mesh's public services behind the derived filter for weeks, and a new line in its ruleset is a
// change to a production firewall nobody asked for.
func TestAMachineThatNamesNoNetworksIsFilteredAsBefore(t *testing.T) {
rules := []Rule{{Port: 443, Protocol: "tcp", From: FromEverywhere, Because: []string{"proxy"}}}
said := AsNftables(rules, []string{"10.10.0.1"}, true, []int{4222}, nil)
quiet := AsNftables(rules, []string{"10.10.0.1"}, true, []int{4222}, []string{})
if said != quiet {
t.Errorf("nil and empty render differently:\n%s\n---\n%s", said, quiet)
}
if strings.Contains(said, "a network this machine routes") {
t.Errorf("a machine that named nothing carries a line about what it routes:\n%s", said)
}
}
// The two families are matched differently, and one set holding both is a syntax error — a ruleset
// that does not load is a machine filtering nothing while its unit reports a fault.
func TestARoutedNetworkIsMatchedInItsOwnFamily(t *testing.T) {
ruleset := AsNftables(nil, nil, false, nil, []string{"10.0.0.0/8", "fd00::/8"})
if !strings.Contains(ruleset, "ip saddr 10.0.0.0/8 accept") {
t.Errorf("the v4 network is not matched as ip saddr:\n%s", ruleset)
}
if !strings.Contains(ruleset, "ip6 saddr fd00::/8 accept") {
t.Errorf("the v6 network is not matched as ip6 saddr:\n%s", ruleset)
}
if strings.Contains(ruleset, "ip saddr fd00::/8") {
t.Errorf("a v6 network is matched as ip saddr, which nftables refuses:\n%s", ruleset)
}
}
// chainOf is one chain's body, so a test about the forward chain cannot pass on a line in the input
// chain that happens to look the same.
func chainOf(t *testing.T, ruleset, name string) string {
t.Helper()
start := strings.Index(ruleset, "chain "+name+" {")
if start < 0 {
t.Fatalf("no chain %q in:\n%s", name, ruleset)
}
rest := ruleset[start:]
end := strings.Index(rest, "\n\t}")
if end < 0 {
t.Fatalf("chain %q does not end:\n%s", name, rest)
}
return rest[:end]
}
+14 -14
View File
@@ -79,7 +79,7 @@ func TestTwoModulesWantingOnePortAreBothNamed(t *testing.T) {
t.Fatalf("a module that wanted this port open is not named: %+v", rules[0]) t.Fatalf("a module that wanted this port open is not named: %+v", rules[0])
} }
// The consequence, which is the reason this matters: removing web must not read as closing 443. // The consequence, which is the reason this matters: removing web must not read as closing 443.
nft := AsNftables(rules, nil, false, nil) nft := AsNftables(rules, nil, false, nil, nil)
if !strings.Contains(nft, "web") || !strings.Contains(nft, "board") { if !strings.Contains(nft, "web") || !strings.Contains(nft, "board") {
t.Fatalf("the rendered rule set does not name both sources:\n%s", nft) t.Fatalf("the rendered rule set does not name both sources:\n%s", nft)
} }
@@ -107,7 +107,7 @@ func TestAPortOpenToEveryoneIsNotAlsoRestrictedToTheMesh(t *testing.T) {
func TestWhatNoModuleDeclaredIsClosed(t *testing.T) { func TestWhatNoModuleDeclaredIsClosed(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "web", Listens: []Listening{{Port: 443, From: FromEverywhere}}}, {Module: "web", Listens: []Listening{{Port: 443, From: FromEverywhere}}},
}}, nil), []string{"198.51.100.2"}, false, nil) }}, nil), []string{"198.51.100.2"}, false, nil, nil)
// Naming the chain, not just the policy: the forward chain drops too, and an assertion on // Naming the chain, not just the policy: the forward chain drops too, and an assertion on
// "policy drop" alone passes while the input chain accepts everything. It did, once, here. // "policy drop" alone passes while the input chain accepts everything. It did, once, here.
if !strings.Contains(nft, "type filter hook input priority filter; policy drop;") { if !strings.Contains(nft, "type filter hook input priority filter; policy drop;") {
@@ -134,7 +134,7 @@ func TestWhatNoModuleDeclaredIsClosed(t *testing.T) {
// `flush ruleset` would do the first and not the second: it empties every table on the machine, // `flush ruleset` would do the first and not the second: it empties every table on the machine,
// including the ones the container runtime writes for its bridges. // including the ones the container runtime writes for its bridges.
func TestReloadingReplacesOnlyTheMeshsOwnRules(t *testing.T) { func TestReloadingReplacesOnlyTheMeshsOwnRules(t *testing.T) {
nft := AsNftables(nil, nil, false, nil) nft := AsNftables(nil, nil, false, nil, nil)
if strings.Contains(nft, "flush ruleset") { if strings.Contains(nft, "flush ruleset") {
t.Fatalf("loading the rule set empties every table on the machine:\n%s", nft) t.Fatalf("loading the rule set empties every table on the machine:\n%s", nft)
} }
@@ -160,7 +160,7 @@ func TestReloadingReplacesOnlyTheMeshsOwnRules(t *testing.T) {
// So the chain exists and denies by default, and the runtime's own networks are allowed explicitly // So the chain exists and denies by default, and the runtime's own networks are allowed explicitly
// — which is how the system being replaced has been doing it on these machines for months. // — which is how the system being replaced has been doing it on these machines for months.
func TestWhatIsForwardedIsGovernedToo(t *testing.T) { func TestWhatIsForwardedIsGovernedToo(t *testing.T) {
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil) nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, nil)
if !strings.Contains(nft, "hook forward priority filter; policy drop") { if !strings.Contains(nft, "hook forward priority filter; policy drop") {
t.Fatalf("forwarded traffic is not governed, so container ports are open:\n%s", nft) t.Fatalf("forwarded traffic is not governed, so container ports are open:\n%s", nft)
} }
@@ -168,7 +168,7 @@ func TestWhatIsForwardedIsGovernedToo(t *testing.T) {
// And containers keep working, which is the whole reason the chain was left out before. // And containers keep working, which is the whole reason the chain was left out before.
func TestTheRuntimesOwnNetworksKeepWorking(t *testing.T) { func TestTheRuntimesOwnNetworksKeepWorking(t *testing.T) {
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil) nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, nil)
for _, network := range []string{"172.16.0.0/12", "192.168.128.0/17"} { for _, network := range []string{"172.16.0.0/12", "192.168.128.0/17"} {
if !strings.Contains(nft, "ip saddr "+network+" accept") { if !strings.Contains(nft, "ip saddr "+network+" accept") {
t.Fatalf("%s is not allowed, so denying by default stops every container:\n%s", network, nft) t.Fatalf("%s is not allowed, so denying by default stops every container:\n%s", network, nft)
@@ -183,7 +183,7 @@ func TestTheRuntimesOwnNetworksKeepWorking(t *testing.T) {
func TestAPublishedPortIsMatchedByWhatWasAskedFor(t *testing.T) { func TestAPublishedPortIsMatchedByWhatWasAskedFor(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "web", Listens: []Listening{{Port: 8080, From: FromEverywhere}}}, {Module: "web", Listens: []Listening{{Port: 8080, From: FromEverywhere}}},
}}, nil), []string{"198.51.100.2"}, false, nil) }}, nil), []string{"198.51.100.2"}, false, nil, nil)
if !strings.Contains(nft, "ct original proto-dst 8080 accept") { if !strings.Contains(nft, "ct original proto-dst 8080 accept") {
t.Fatalf("the forwarded rule does not match the port a client asked for:\n%s", nft) t.Fatalf("the forwarded rule does not match the port a client asked for:\n%s", nft)
} }
@@ -193,7 +193,7 @@ func TestAPublishedPortIsMatchedByWhatWasAskedFor(t *testing.T) {
func TestAMeshScopedPortIsMeshScopedWhenForwarded(t *testing.T) { func TestAMeshScopedPortIsMeshScopedWhenForwarded(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}}, {Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}},
}}, nil), []string{"198.51.100.2"}, false, nil) }}, nil), []string{"198.51.100.2"}, false, nil, nil)
if !strings.Contains(nft, "ip saddr { 198.51.100.2 } ct original proto-dst 5432 accept") { if !strings.Contains(nft, "ip saddr { 198.51.100.2 } ct original proto-dst 5432 accept") {
t.Fatalf("a mesh-only port is reachable from anywhere once forwarded:\n%s", nft) t.Fatalf("a mesh-only port is reachable from anywhere once forwarded:\n%s", nft)
} }
@@ -203,7 +203,7 @@ func TestAMeshScopedPortIsMeshScopedWhenForwarded(t *testing.T) {
func TestFromTheMeshIsTheNodesTheMeshKnows(t *testing.T) { func TestFromTheMeshIsTheNodesTheMeshKnows(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}}, {Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}},
}}, nil), []string{"198.51.100.2", "198.51.100.3"}, false, nil) }}, nil), []string{"198.51.100.2", "198.51.100.3"}, false, nil, nil)
if !strings.Contains(nft, "ip saddr { 198.51.100.2, 198.51.100.3 } tcp dport 5432 accept") { if !strings.Contains(nft, "ip saddr { 198.51.100.2, 198.51.100.3 } tcp dport 5432 accept") {
t.Fatalf("a mesh-scoped port was not restricted to the mesh's addresses:\n%s", nft) t.Fatalf("a mesh-scoped port was not restricted to the mesh's addresses:\n%s", nft)
} }
@@ -213,7 +213,7 @@ func TestFromTheMeshIsTheNodesTheMeshKnows(t *testing.T) {
func TestAMeshPortOnANodeWithNoMeshIsClosedAndSaysSo(t *testing.T) { func TestAMeshPortOnANodeWithNoMeshIsClosedAndSaysSo(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}}, {Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}},
}}, nil), nil, false, nil) }}, nil), nil, false, nil, nil)
if strings.Contains(nft, "dport 5432 accept") { if strings.Contains(nft, "dport 5432 accept") {
t.Fatalf("a port meant for the mesh was opened to everything:\n%s", nft) t.Fatalf("a port meant for the mesh was opened to everything:\n%s", nft)
} }
@@ -226,7 +226,7 @@ func TestAMeshPortOnANodeWithNoMeshIsClosedAndSaysSo(t *testing.T) {
func TestAMachineScopedPortIsNotOpened(t *testing.T) { func TestAMachineScopedPortIsNotOpened(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "cache", Listens: []Listening{{Port: 6379, From: FromMachine}}}, {Module: "cache", Listens: []Listening{{Port: 6379, From: FromMachine}}},
}}, nil), []string{"198.51.100.2"}, false, nil) }}, nil), []string{"198.51.100.2"}, false, nil, nil)
if strings.Contains(nft, "dport 6379 accept") { if strings.Contains(nft, "dport 6379 accept") {
t.Fatalf("a port for this machine only was opened to the network:\n%s", nft) t.Fatalf("a port for this machine only was opened to the network:\n%s", nft)
} }
@@ -268,7 +268,7 @@ func TestAskingForTheRuleSetWithNowhereToPutItIsRefused(t *testing.T) {
func TestAMeshOnBothAddressFamiliesRendersBoth(t *testing.T) { func TestAMeshOnBothAddressFamiliesRendersBoth(t *testing.T) {
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{ nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
{Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}}, {Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}},
}}, nil), []string{"198.51.100.2", "2001:db8::2"}, false, nil) }}, nil), []string{"198.51.100.2", "2001:db8::2"}, false, nil, nil)
if !strings.Contains(nft, "ip saddr { 198.51.100.2 } tcp dport 5432 accept") { if !strings.Contains(nft, "ip saddr { 198.51.100.2 } tcp dport 5432 accept") {
t.Fatalf("the machines with v4 addresses were dropped:\n%s", nft) t.Fatalf("the machines with v4 addresses were dropped:\n%s", nft)
} }
@@ -672,7 +672,7 @@ func TestExposureRefusesAPortNotListenedOnAndABadSource(t *testing.T) {
// loading the rules lives on conntrack until it drops, and then the machine is reached from a // loading the rules lives on conntrack until it drops, and then the machine is reached from a
// rescue console (novox/hq issue 047). // rescue console (novox/hq issue 047).
func TestSSHIsOpenFromTheMeshEvenWhenNothingIsAssigned(t *testing.T) { func TestSSHIsOpenFromTheMeshEvenWhenNothingIsAssigned(t *testing.T) {
nft := AsNftables(nil, []string{"198.51.100.2", "198.51.100.3"}, false, nil) nft := AsNftables(nil, []string{"198.51.100.2", "198.51.100.3"}, false, nil, nil)
if !strings.Contains(nft, "ip saddr { 198.51.100.2, 198.51.100.3 } tcp dport 22 accept") { if !strings.Contains(nft, "ip saddr { 198.51.100.2, 198.51.100.3 } tcp dport 22 accept") {
t.Fatalf("ssh is not open to the mesh, so a machine can lock everyone out:\n%s", nft) t.Fatalf("ssh is not open to the mesh, so a machine can lock everyone out:\n%s", nft)
} }
@@ -685,7 +685,7 @@ func TestSSHIsOpenFromTheMeshEvenWhenNothingIsAssigned(t *testing.T) {
// And from outside as well, on a machine that faces outward — because that is the way in when the // And from outside as well, on a machine that faces outward — because that is the way in when the
// private network is the thing that broke. // private network is the thing that broke.
func TestSSHIsOpenFromOutsideOnAMachineThatFacesIt(t *testing.T) { func TestSSHIsOpenFromOutsideOnAMachineThatFacesIt(t *testing.T) {
nft := AsNftables(nil, []string{"198.51.100.2"}, true, nil) nft := AsNftables(nil, []string{"198.51.100.2"}, true, nil, nil)
if !strings.Contains(nft, "\t\ttcp dport 22 accept") { if !strings.Contains(nft, "\t\ttcp dport 22 accept") {
t.Fatalf("a machine reachable from outside does not answer ssh there:\n%s", nft) t.Fatalf("a machine reachable from outside does not answer ssh there:\n%s", nft)
} }
@@ -697,7 +697,7 @@ func TestSSHIsOpenFromOutsideOnAMachineThatFacesIt(t *testing.T) {
// to narrow the rule to, so narrowing it shuts the port entirely — on the first machine anybody // to narrow the rule to, so narrowing it shuts the port entirely — on the first machine anybody
// adopts, reached over the network, closed by the act of adopting it. // adopts, reached over the network, closed by the act of adopting it.
func TestSSHIsNeverLeftWithoutARule(t *testing.T) { func TestSSHIsNeverLeftWithoutARule(t *testing.T) {
nft := AsNftables(nil, nil, false, nil) nft := AsNftables(nil, nil, false, nil, nil)
if !strings.Contains(nft, "tcp dport 22 accept") { if !strings.Contains(nft, "tcp dport 22 accept") {
t.Fatalf("a machine with no mesh addresses has no ssh rule, so adopting it locks it:\n%s", nft) t.Fatalf("a machine with no mesh addresses has no ssh rule, so adopting it locks it:\n%s", nft)
} }
+24
View File
@@ -225,6 +225,19 @@ type Manifest struct {
// subscription to the queue it writes to (design 29 §2). // subscription to the queue it writes to (design 29 §2).
Uses []string `json:"uses,omitempty"` Uses []string `json:"uses,omitempty"`
// Prepares says this module has state that must be brought to the shape this version needs
// before this version runs, and that the module's own code does it (novox/hq ADR 0135).
//
// **A word, not an arrangement.** The mesh runs the module's own program in its preparation
// mode, in the module's own context — every binding, credential and setting its code receives,
// because it *is* its code. Nothing here names a container, a command, a mount or a variable:
// the module already said all of that once, and a second copy is a second thing to drift.
//
// **Declared, never inferred.** The control plane cannot read what is inside an artifact, so a
// module that ships a migration and does not say this breaks on its first upgrade. That is
// stated in the record rather than guarded here, because nothing mechanical can guard it.
Prepares bool `json:"prepares,omitempty"`
// Tools are the tools this module answers — request and reply, awaited. // Tools are the tools this module answers — request and reply, awaited.
// //
// **New, and not `serves`**, which this manifest already uses for the facts a consumer needs // **New, and not `serves`**, which this manifest already uses for the facts a consumer needs
@@ -1276,6 +1289,17 @@ func ParseManifest(raw []byte) (Manifest, error) {
"program that reads what the mesh delivered and reconciles", "program that reads what the mesh delivered and reconciles",
m.Module, r["id"])) m.Module, r["id"]))
} }
// **A module that prepares its state must have code the mesh can run** (novox/hq ADR 0135). The
// preparation is the module's own program in its preparation mode, so it is derived from the
// resource that runs that program — and a module declaring none has asked for something the mesh
// cannot compose. Said here, where the manifest is read, rather than by a declaration that
// quietly prepares nothing.
if m.Prepares && preparationTarget(m) == "" {
problems = append(problems, fmt.Sprintf(
"%s says it prepares its state, and declares no container running an artifact it built — "+
"the preparation is this module's own program, so there has to be one for the mesh to "+
"run it in", m.Module))
}
// **A run-once container is a step the host runs to completion** (novox/hq ADR 0052). It is a // **A run-once container is a step the host runs to completion** (novox/hq ADR 0052). It is a
// boolean modifier on the container shape — the host runs the container, requires it to exit 0, // boolean modifier on the container shape — the host runs the container, requires it to exit 0,
// and starts whatever the declaration places after it only once it has. A value that is not a // and starts whatever the declaration places after it only once it has. A value that is not a
+141
View File
@@ -0,0 +1,141 @@
package catalogue
import (
"encoding/json"
"fmt"
"strings"
"testing"
)
// A module version prepares its state before it runs (novox/hq ADR 0135).
//
// What the mesh derives is the module's own resource, run once with one word, placed immediately in
// front of the thing it prepares for. What matters in these tests is that the derivation is a copy
// rather than a second description: the failure it replaces was a hand-written step repeating six
// fields of the resource it preceded, each free to drift from it.
func aPreparingModule() Manifest {
return Manifest{
Module: "gitea",
Prepares: true,
Resources: []map[string]any{
{"id": "state", "type": "directory", "path": "/var/lib/gitea", "mode": "0700"},
{"id": "server", "type": "container", "name": "mesh-gitea-server",
"image": "gitea/gitea@" + digest, "ports": []any{"3000:3000"}},
{"id": "runtime", "type": "container", "name": "mesh-gitea",
"image": ArtifactStoreScheme + "gitea/runtime@" + digest, "network": "host",
"env": map[string]any{"MESH_GITEA_STATE_DIR": "/run/state"},
"volumes": []any{"/var/lib/gitea:/run/state:ro"},
"ports": []any{"9000:9000"}},
},
}
}
func declaredFor(t *testing.T, m Manifest) []map[string]any {
t.Helper()
// The manifest as the mesh holds it: a build resolved the module's own artifact into the
// reference it recorded, which is also how the composition knows whose code a resource runs.
out, err := Resolution{Node: "anchor", Modules: []Manifest{m}}.Declaration(
Rendering{ArtifactStore: "anchor.internal:5100"})
if err != nil {
t.Fatal(err)
}
return out
}
func idsOf(resources []map[string]any) []string {
var ids []string
for _, r := range resources {
ids = append(ids, fmt.Sprint(r["id"]))
}
return ids
}
// The step runs the module's own code, and comes immediately before it — not before the upstream
// server the module packages, which may be the very thing the state lives in.
func TestThePreparationRunsTheModulesOwnCodeAndComesRightBeforeIt(t *testing.T) {
out := declaredFor(t, aPreparingModule())
ids := idsOf(out)
at := -1
for i, id := range ids {
if id == "gitea.runtime-prepare" {
at = i
}
}
if at < 0 {
t.Fatalf("nothing prepares this module's state: %v", ids)
}
// A module's name may contain a dot, so a resource's module is everything before the last one —
// which the derived id must not add to, or a machine reads the wrong owner from it.
if strings.Count("gitea.runtime-prepare", ".") != 1 {
t.Fatal("the derived id adds a dot, so what owns it cannot be read from it")
}
if ids[at+1] != "gitea.runtime" {
t.Fatalf("the preparation is not immediately before the module's own code: %v", ids)
}
for _, id := range ids[:at] {
if id == "gitea.runtime" {
t.Fatalf("the module's own code runs before its state is prepared: %v", ids)
}
}
}
// It is given exactly what the module's own code is given. Asserted field by field against the
// resource it was derived from, because writing it twice is the fault this replaces.
func TestThePreparationIsGivenWhatTheModuleIsGiven(t *testing.T) {
out := declaredFor(t, aPreparingModule())
declared := byID(out)
step, workload := declared["gitea.runtime-prepare"], declared["gitea.runtime"]
if step == nil || workload == nil {
t.Fatalf("expected both, got %v", idsOf(out))
}
for _, field := range []string{"image", "network", "env", "volumes", "type"} {
if fmt.Sprint(step[field]) != fmt.Sprint(workload[field]) {
t.Errorf("the preparation's %s is %v and the module's is %v", field, step[field], workload[field])
}
}
if once, _ := step["run-once"].(bool); !once {
t.Error("the preparation is not a step, so nothing waits for it and nothing is gated by it")
}
if fmt.Sprint(step["args"]) != fmt.Sprint([]any{PreparationArgument}) {
t.Errorf("the preparation is asked for as %v", step["args"])
}
if fmt.Sprint(step["name"]) == fmt.Sprint(workload["name"]) {
t.Error("the preparation and the workload have one name, so one removes the other")
}
// A published port cannot be bound twice, and the version being replaced is still running.
if _, published := step["ports"]; published {
t.Errorf("the preparation publishes a port the running version holds: %v", step["ports"])
}
}
// A module that says nothing about preparing gets nothing, which is most modules.
func TestAModuleThatPreparesNothingGetsNoStep(t *testing.T) {
m := aPreparingModule()
m.Prepares = false
for _, id := range idsOf(declaredFor(t, m)) {
if id == "gitea.runtime-prepare" {
t.Fatal("a module that prepares nothing was given a preparation")
}
}
}
// A module whose own code the mesh cannot find has nothing to ask, and saying so where the manifest
// is read beats a declaration that quietly prepares nothing.
func TestAModuleThatPreparesAndRunsNoneOfItsOwnCodeIsRefused(t *testing.T) {
m := Manifest{
Module: "gitea",
Prepares: true,
Resources: []map[string]any{
// Only the upstream server it packages: nothing here runs gitea's own code.
{"id": "server", "type": "container", "name": "mesh-gitea-server", "image": "gitea/gitea@" + digest},
},
}
raw, err := json.Marshal(m)
if err != nil {
t.Fatal(err)
}
if _, err := ParseManifest(raw); err == nil {
t.Fatal("a module that prepares its state with nothing of its own to run was accepted")
}
}
+1 -1
View File
@@ -12,5 +12,5 @@ func TestPrintRehearsalRuleset(t *testing.T) {
rules := mustFilter(t, Resolution{Modules: []Manifest{ rules := mustFilter(t, Resolution{Modules: []Manifest{
{Module: "pub", Listens: []Listening{{Port: 8099, From: FromMesh, Why: "the thing it serves"}}}, {Module: "pub", Listens: []Listening{{Port: 8099, From: FromMesh, Why: "the thing it serves"}}},
}}, nil) }}, nil)
t.Log("\n" + AsNftables(rules, []string{"192.0.2.20"}, true, nil)) t.Log("\n" + AsNftables(rules, []string{"192.0.2.20"}, true, nil, nil))
} }
+5 -1
View File
@@ -48,7 +48,11 @@ type Seat struct {
// //
// In the order a person reads it: the mesh's own, then a node's. // In the order a person reads it: the mesh's own, then a node's.
var defaultSeats = []Seat{ var defaultSeats = []Seat{
{Name: "mesh-controller", Scope: ScopeMesh, Decision: "novox/hq ADR 0079"}, // The control plane states what it did under the seat it holds (novox/hq ADR 0134): a role's
// events belong to the role, so they keep their address while the holder is replaced. No accepts,
// so no work queue is raised for it — only what its holder may say.
{Name: "mesh-controller", Scope: ScopeMesh, Decision: "novox/hq ADR 0079",
Emits: []string{"applied", "refused", "built-before"}},
{Name: "mesh-store", Scope: ScopeMesh, Delivers: "postgres-database", Decision: "novox/hq ADR 0079"}, {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 // **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 // ADR 0127 separated them: `amqp` is a backing service a module may require, and this seat is
+12 -12
View File
@@ -30,12 +30,12 @@ func TestRefusedAndFailedAreDifferentSituations(t *testing.T) {
refuser := nodeNamed(t, inv, "refuser") refuser := nodeNamed(t, inv, "refuser")
failer := nodeNamed(t, inv, "failer") failer := nodeNamed(t, inv, "failer")
if err := inv.RecordDoing(ctx, refuser, Doing{ if _, err := inv.RecordDoing(ctx, refuser, Doing{
Outcome: OutcomeRefused, Refused: "resource \"x\": a file needs a path", Outcome: OutcomeRefused, Refused: "resource \"x\": a file needs a path",
}); err != nil { }); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := inv.RecordDoing(ctx, failer, Doing{ if _, err := inv.RecordDoing(ctx, failer, Doing{
Outcome: OutcomeFailed, Outcome: OutcomeFailed,
Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}},
Applied: 4, Applied: 4,
@@ -70,7 +70,7 @@ func TestAMachineDoingWhatItWasToldIsNotOnTheList(t *testing.T) {
inv := fresh(t) inv := fresh(t)
ctx := context.Background() ctx := context.Background()
id := nodeNamed(t, inv, "fine") id := nodeNamed(t, inv, "fine")
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 6}); err != nil { if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 6}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
wrong, err := inv.NotDoingWhatTheyWereTold(ctx) wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
@@ -97,12 +97,12 @@ func TestTheLastReportReplacesTheOneBefore(t *testing.T) {
inv := fresh(t) inv := fresh(t)
ctx := context.Background() ctx := context.Background()
id := nodeNamed(t, inv, "recovered") id := nodeNamed(t, inv, "recovered")
if err := inv.RecordDoing(ctx, id, Doing{ if _, err := inv.RecordDoing(ctx, id, Doing{
Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "a", Error: "no"}}, Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "a", Error: "no"}},
}); err != nil { }); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 3}); err != nil { if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 3}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
wrong, err := inv.NotDoingWhatTheyWereTold(ctx) wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
@@ -141,7 +141,7 @@ func TestWhatANodeSaidGoesWhenTheNodeDoes(t *testing.T) {
inv := fresh(t) inv := fresh(t)
ctx := context.Background() ctx := context.Background()
id := nodeNamed(t, inv, "leaving") id := nodeNamed(t, inv, "leaving")
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeFailed}); err != nil { if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeFailed}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil { if _, err := inv.store.Pool().Exec(ctx, `delete from node where name = 'leaving'`); err != nil {
@@ -257,7 +257,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
id := nodeNamed(t, inv, "looping") id := nodeNamed(t, inv, "looping")
same := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image"}}} same := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image"}}}
if err := inv.RecordDoing(ctx, id, same); err != nil { if _, err := inv.RecordDoing(ctx, id, same); err != nil {
t.Fatal(err) t.Fatal(err)
} }
first, _, err := inv.DoingOf(ctx, "looping") first, _, err := inv.DoingOf(ctx, "looping")
@@ -269,7 +269,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
} }
for range StuckAfter - 1 { for range StuckAfter - 1 {
if err := inv.RecordDoing(ctx, id, same); err != nil { if _, err := inv.RecordDoing(ctx, id, same); err != nil {
t.Fatal(err) t.Fatal(err)
} }
} }
@@ -287,7 +287,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
// The same resource failing with different words — a duration, a counter — is still the same // The same resource failing with different words — a duration, a counter — is still the same
// failure: it is the resource that loops, not the sentence. // failure: it is the resource that loops, not the sentence.
reworded := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image (after 31s)"}}} reworded := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "img", Error: "no such image (after 31s)"}}}
if err := inv.RecordDoing(ctx, id, reworded); err != nil { if _, err := inv.RecordDoing(ctx, id, reworded); err != nil {
t.Fatal(err) t.Fatal(err)
} }
still, _, err := inv.DoingOf(ctx, "looping") still, _, err := inv.DoingOf(ctx, "looping")
@@ -300,7 +300,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
// A different failure is a new situation, not a longer one. // A different failure is a new situation, not a longer one.
other := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}} other := Doing{Outcome: OutcomeFailed, Failed: []FailedResource{{ID: "svc", Error: "unit not found"}}}
if err := inv.RecordDoing(ctx, id, other); err != nil { if _, err := inv.RecordDoing(ctx, id, other); err != nil {
t.Fatal(err) t.Fatal(err)
} }
changed, _, err := inv.DoingOf(ctx, "looping") changed, _, err := inv.DoingOf(ctx, "looping")
@@ -312,7 +312,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
} }
// And a clean apply clears it: the machine is doing what it was told, since nothing. // And a clean apply clears it: the machine is doing what it was told, since nothing.
if err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 2}); err != nil { if _, err := inv.RecordDoing(ctx, id, Doing{Outcome: OutcomeApplied, Applied: 2}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
fine, _, err := inv.DoingOf(ctx, "looping") fine, _, err := inv.DoingOf(ctx, "looping")
@@ -324,7 +324,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
} }
// The list of what is wrong carries the count, so `status` can say it. // The list of what is wrong carries the count, so `status` can say it.
if err := inv.RecordDoing(ctx, id, same); err != nil { if _, err := inv.RecordDoing(ctx, id, same); err != nil {
t.Fatal(err) t.Fatal(err)
} }
wrong, err := inv.NotDoingWhatTheyWereTold(ctx) wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
@@ -0,0 +1,19 @@
-- The networks a machine routes for what it hosts, beyond the container runtime's own defaults.
--
-- novox/hq ADR 0137. The derived packet filter denies forwarding by default and then allows the
-- container runtime's two default pools, named in the controller's code with a comment saying that
-- a machine configured otherwise "needs this to say so" — and no way to say it. So the filter was
-- correct only on a machine whose runtime used the defaults, and silently wrong on any other.
--
-- Measured on 2026-09-28: flipping a workstation to the derived filter cut egress for five of its
-- container networks and for every network its test beds create, because those are allocated from
-- ranges the two defaults do not cover. Nothing reported a fault; the containers simply could not
-- reach anything.
--
-- A node-level fact, beside the node's public domain and for the same reason: it is a property of
-- the machine, not of whichever module happens to load the filter today. Swapping that module must
-- not lose it.
--
-- Null for a machine that routes nothing but the runtime's defaults, which is the ordinary case and
-- what every machine held before this column existed.
alter table node add column routed_networks jsonb;
+93 -18
View File
@@ -9,6 +9,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"net"
"sort" "sort"
"strings" "strings"
"time" "time"
@@ -518,6 +519,69 @@ func (i *Inventory) PublicDomainOf(ctx context.Context, name string) (string, er
return *domain, nil return *domain, nil
} }
// SetRoutedNetworks records the networks this machine routes for what it hosts, beyond the
// container runtime's own default pools.
//
// A node-level fact (novox/hq ADR 0137), beside the node's public domain: the machine routes them,
// not whichever module loads the filter, so swapping that module must not lose them. Added to the
// runtime's defaults rather than replacing them, so a machine that says one range does not lose the
// ranges its containers were already using. An empty list clears it.
//
// Each entry is checked as a CIDR here rather than at render time: an address that does not parse
// becomes a line nftables refuses, and a refused ruleset is a machine that filters nothing while
// its service reports a configuration fault.
func (i *Inventory) SetRoutedNetworks(ctx context.Context, name string, networks []string) error {
node, err := i.NodeByName(ctx, name)
if err != nil {
return err
}
var kept []string
for _, n := range networks {
n = strings.TrimSpace(n)
if n == "" {
continue
}
if _, _, err := net.ParseCIDR(n); err != nil {
return fmt.Errorf("%q is not a network in CIDR form (10.0.0.0/8, 192.168.0.0/16): %w",
n, err)
}
kept = append(kept, n)
}
if len(kept) == 0 {
_, err = i.store.Pool().Exec(ctx,
`update node set routed_networks = null where id = $1`, node.ID)
return err
}
body, err := json.Marshal(kept)
if err != nil {
return err
}
_, err = i.store.Pool().Exec(ctx,
`update node set routed_networks = $2 where id = $1`, node.ID, string(body))
return err
}
// RoutedNetworksOf is the networks a machine routes for what it hosts, empty when it has named none.
func (i *Inventory) RoutedNetworksOf(ctx context.Context, name string) ([]string, error) {
var body []byte
err := i.store.Pool().QueryRow(ctx,
`select routed_networks from node where name = $1`, name).Scan(&body)
if errors.Is(err, pgx.ErrNoRows) {
return nil, fmt.Errorf("%w: %s", ErrNoSuchNode, name)
}
if err != nil {
return nil, err
}
if len(body) == 0 {
return nil, nil
}
var networks []string
if err := json.Unmarshal(body, &networks); err != nil {
return nil, fmt.Errorf("the networks recorded for %s are not a list: %w", name, err)
}
return networks, nil
}
// RecordOverlayKey keeps the public half a node generated. // RecordOverlayKey keeps the public half a node generated.
func (i *Inventory) RecordOverlayKey(ctx context.Context, node, key string) error { func (i *Inventory) RecordOverlayKey(ctx context.Context, node, key string) error {
if strings.TrimSpace(key) == "" { if strings.TrimSpace(key) == "" {
@@ -670,31 +734,39 @@ func sameFailure(a, b Doing) bool {
// a clean apply clears both (novox/hq 04-ISSUES/065). The previous row is read first and the // a clean apply clears both (novox/hq 04-ISSUES/065). The previous row is read first and the
// comparison made here, so "the same" is a rule this package states rather than a jsonb equality // comparison made here, so "the same" is a rule this package states rather than a jsonb equality
// that would restart the count on a changed word in an error. // that would restart the count on a changed word in an error.
func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error { // **And whether this report was news**, which is what makes a fact about it worth stating (novox/hq
// ADR 0134). A machine reconciles continuously and reports each time; the same outcome about the same
// declaration is the same state said again, and a fact per report would be a fact per minute per
// machine that tells nobody anything. Read here because the previous row is read here anyway.
func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) (news bool, err error) {
failed, err := json.Marshal(d.Failed) failed, err := json.Marshal(d.Failed)
if err != nil { if err != nil {
return err return false, err
}
var before Doing
var beforeFailed []byte
found := i.store.Pool().QueryRow(ctx,
`select outcome, refused, failed, failing_since, failures, coalesce(declared,'')
from node_report where node = $1`,
node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times,
&before.Declared)
switch {
case errors.Is(found, pgx.ErrNoRows):
news = true
case found != nil:
return false, found
default:
if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil {
return false, err
}
news = before.Outcome != d.Outcome || before.Declared != d.Declared || !sameFailure(before, d)
} }
var since *time.Time var since *time.Time
times := 0 times := 0
if d.Outcome != OutcomeApplied { if d.Outcome != OutcomeApplied {
var before Doing
var beforeFailed []byte
err := i.store.Pool().QueryRow(ctx,
`select outcome, refused, failed, failing_since, failures from node_report where node = $1`,
node).Scan(&before.Outcome, &before.Refused, &beforeFailed, &before.Since, &before.Times)
switch {
case errors.Is(err, pgx.ErrNoRows):
case err != nil:
return err
default:
if err := json.Unmarshal(beforeFailed, &before.Failed); err != nil {
return err
}
}
now := time.Now() now := time.Now()
since, times = &now, 1 since, times = &now, 1
if err == nil && sameFailure(before, d) && before.Since != nil { if found == nil && sameFailure(before, d) && before.Since != nil {
since, times = before.Since, before.Times+1 since, times = before.Since, before.Times+1
} }
} }
@@ -707,7 +779,10 @@ func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error
declared = excluded.declared, declared = excluded.declared,
failing_since = excluded.failing_since, failures = excluded.failures`, failing_since = excluded.failing_since, failures = excluded.failures`,
node, d.Outcome, d.Refused, failed, d.Applied, d.Declared, since, times) node, d.Outcome, d.Refused, failed, d.Applied, d.Declared, since, times)
return err if err != nil {
return false, err
}
return news, nil
} }
// NotDoingWhatTheyWereTold is every machine whose last report was not a clean apply. // NotDoingWhatTheyWereTold is every machine whose last report was not a clean apply.
+33
View File
@@ -30,6 +30,11 @@ type Bus interface {
// (design 29 §4, the *state* shape). // (design 29 §4, the *state* shape).
PublishDeclaration(ctx context.Context, node string, body []byte) error PublishDeclaration(ctx context.Context, node string, body []byte) error
// PublishSeatEvent states a fact under a role's own name, for the holder of that role. A
// module's event is addressed to the module; a role's is addressed to the role, so it keeps
// meaning when the holder changes (novox/hq ADR 0121, ADR 0129).
PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error
// AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves // AskTool sends one question to a module's tool and awaits one answer. A tool nobody serves
// must say so **at once** rather than after the whole wait: the difference between "that // must say so **at once** rather than after the whole wait: the difference between "that
// module is down" and "that tool is slow" is the first thing a person asking wants. // module is down" and "that tool is slow" is the first thing a person asking wants.
@@ -81,6 +86,13 @@ func EventSubject(source, key string) string {
return "mesh.mod." + source + ".event." + key return "mesh.mod." + source + ".event." + key
} }
// SeatEventSubject is where a role's own event lands. Derived from the role, never from its holder:
// a fact about the build machine or about the control plane keeps its address when the module holding
// that role is replaced (novox/hq ADR 0121, ADR 0129).
func SeatEventSubject(seat, event string) string {
return "mesh.seat." + seat + ".event." + event
}
// DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so // DeclareSubject is where one node's declaration lands. Last-per-subject on the NODES stream, so
// a node that was away gets exactly the current one and a replayed older one is refused by // a node that was away gets exactly the current one and a replayed older one is refused by
// sequence — the wire-level answer to novox/hq issue 107. // sequence — the wire-level answer to novox/hq issue 107.
@@ -112,6 +124,27 @@ func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, bo
return nil return nil
} }
// PublishSeatEvent states a role's own fact. Same envelope as a module's event and a different
// address: the source header is the role, because that is what the fact is about.
func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error {
id, err := eventID()
if err != nil {
return err
}
h := nats.Header{}
h.Set("x-event-id", id)
h.Set("x-source", seat)
_, err = b.JS.PublishMsg(&nats.Msg{
Subject: SeatEventSubject(seat, event),
Header: h,
Data: body,
}, nats.MsgId(id), nats.Context(ctx))
if err != nil {
return fmt.Errorf("stating %s of the %s seat: %w", event, seat, err)
}
return nil
}
func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error { func (b OverNATS) PublishDeclaration(ctx context.Context, node string, body []byte) error {
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx)) _, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
if err != nil { if err != nil {
+16 -13
View File
@@ -265,7 +265,7 @@ func (e Enrolment) Outstanding(ctx context.Context, node string) (string, error)
return e.Inventory.Outstanding(ctx, node) return e.Inventory.Outstanding(ctx, node)
} }
func (e Enrolment) Heard(ctx context.Context, report Report) (err error) { func (e Enrolment) Heard(ctx context.Context, report Report) (news bool, err error) {
// A store that could not be asked right now is said as such, so the report is kept for // A store that could not be asked right now is said as such, so the report is kept for
// another attempt rather than acknowledged and lost (novox/hq issue 082). // another attempt rather than acknowledged and lost (novox/hq issue 082).
defer func() { defer func() {
@@ -274,11 +274,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
} }
}() }()
if report.Node == "" { if report.Node == "" {
return errors.New("a report named no node") return false, errors.New("a report named no node")
} }
node, err := e.Inventory.NodeByName(ctx, report.Node) node, err := e.Inventory.NodeByName(ctx, report.Node)
if err != nil { if err != nil {
return err return false, err
} }
// What an adopted node holds, which firewall it found, and what is reachable on it (novox/hq // What an adopted node holds, which firewall it found, and what is reachable on it (novox/hq
@@ -298,7 +298,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
Port: r.Port, By: r.By, Published: r.Published, ContainerPort: r.ContainerPort}) Port: r.Port, By: r.By, Published: r.Published, ContainerPort: r.ContainerPort})
} }
if err := e.Inventory.RecordAdoption(ctx, node.ID, held, report.Firewall, reachable); err != nil { if err := e.Inventory.RecordAdoption(ctx, node.ID, held, report.Firewall, reachable); err != nil {
return err return false, err
} }
} }
// What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it. // What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it.
@@ -308,7 +308,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
Peers: report.Tunnel.Peers, State: report.Tunnel.State, Note: report.Tunnel.Note, Peers: report.Tunnel.Peers, State: report.Tunnel.State, Note: report.Tunnel.Note,
Kept: report.Tunnel.Kept, Kept: report.Tunnel.Kept,
}); err != nil { }); err != nil {
return err return false, err
} }
} }
// A node taking a found tunnel's key after enrolment (novox/hq ADR 0105). Verified against the // A node taking a found tunnel's key after enrolment (novox/hq ADR 0105). Verified against the
@@ -317,9 +317,9 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
// does not verify or is stale — a refusal, not "not now", so the node hears why. // does not verify or is stale — a refusal, not "not now", so the node hears why.
if report.Rekey != nil { if report.Rekey != nil {
if err := e.rekey(ctx, node, *report.Rekey); err != nil { if err := e.rekey(ctx, node, *report.Rekey); err != nil {
return err return false, err
} }
return e.Inventory.Seen(ctx, node.ID) return false, e.Inventory.Seen(ctx, node.ID)
} }
// A bare word that a node is there is not an account of what the machine did or holds: it // A bare word that a node is there is not an account of what the machine did or holds: it
@@ -334,7 +334,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
if report.Superseded != "" { if report.Superseded != "" {
log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, report.Superseded) log.Printf("%s set aside declaration %s for the newer %s", report.Node, report.Declared, report.Superseded)
} }
return e.Inventory.Seen(ctx, node.ID) return false, e.Inventory.Seen(ctx, node.ID)
} }
// What it did is kept whichever way it went. Until this, a refusal or a failure moved // What it did is kept whichever way it went. Until this, a refusal or a failure moved
// last_seen and the reason went to a log line, so "which machine is not doing what it was // last_seen and the reason went to a log line, so "which machine is not doing what it was
@@ -361,17 +361,20 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
// on top of it (novox/hq ADR 0038). Kept even when the declaration was refused: what the // on top of it (novox/hq ADR 0038). Kept even when the declaration was refused: what the
// machine carries is true regardless of what it thought of the last thing it was sent. // machine carries is true regardless of what it thought of the last thing it was sent.
if err := e.Inventory.RecordCarried(ctx, report.Node, report.Carried); err != nil { if err := e.Inventory.RecordCarried(ctx, report.Node, report.Carried); err != nil {
return err return false, err
} }
if err := e.Inventory.RecordDoing(ctx, node.ID, doing); err != nil { // **Whether this is news** is the store's answer: it holds the previous report, and a machine
return err // that reconciles every minute says the same thing until something changes (novox/hq ADR 0134).
news, err = e.Inventory.RecordDoing(ctx, node.ID, doing)
if err != nil {
return false, err
} }
// A refusal, a failure, or a bare word that the node is there — none of them is an account of // A refusal, a failure, or a bare word that the node is there — none of them is an account of
// what the machine holds, so each moves last_seen and nothing else. Recording a partial list // what the machine holds, so each moves last_seen and nothing else. Recording a partial list
// as though it were the whole would tell a rebuilding node to remove what it still has. // as though it were the whole would tell a rebuilding node to remove what it still has.
if report.Refused != "" || len(report.Failed) > 0 || report.Applied == nil { if report.Refused != "" || len(report.Failed) > 0 || report.Applied == nil {
return e.Inventory.Seen(ctx, node.ID) return news, e.Inventory.Seen(ctx, node.ID)
} }
return e.Inventory.RecordOwned(ctx, node.ID, report.Applied) return news, e.Inventory.RecordOwned(ctx, node.ID, report.Applied)
} }
+33
View File
@@ -49,6 +49,39 @@ func eventID() (string, error) {
return hex.EncodeToString(raw), nil return hex.EncodeToString(raw), nil
} }
// MeshControllerSeat is the role the control plane holds, and therefore where its own facts live: a
// role's events belong to the role, not to whichever container is holding it today (novox/hq ADR 0121,
// ADR 0129). It is what makes them addressable while the control plane itself is being replaced.
const MeshControllerSeat = "mesh-controller"
// The facts the mesh states about its own work (novox/hq ADR 0134).
const (
// KeyApplied: a machine now runs what it was sent.
KeyApplied = "applied"
// KeyRefused: a machine did not take what it was sent, and why.
KeyRefused = "refused"
// KeyBuiltBefore: a build the mesh already held, for a catalogue that asked what it missed. Not
// `built` — that is the build machine's, said as it happens, and a replay is neither.
KeyBuiltBefore = "built-before"
)
// Applied is what a machine now runs, as the mesh states it.
type Applied struct {
Node string `json:"node"`
Declared string `json:"declared,omitempty"`
// Resources is how many the machine applied, not which: the list is the machine's own account
// of itself and belongs in the records, not in a fact every listener has to read past.
Resources int `json:"resources"`
}
// Refused is a machine that would not take what it was sent.
type Refused struct {
Node string `json:"node"`
Declared string `json:"declared,omitempty"`
Refused string `json:"refused,omitempty"`
Failed map[string]string `json:"failed,omitempty"`
}
// KeyModuleBuilt is what the builder announces when it has built something. The catalogue places // KeyModuleBuilt is what the builder announces when it has built something. The catalogue places
// it in the module graph; nothing else need care. // it in the module graph; nothing else need care.
const KeyModuleBuilt = "module.builder.built" const KeyModuleBuilt = "module.builder.built"
+6 -6
View File
@@ -22,7 +22,7 @@ func heardFrom(t *testing.T, report link.Report) (*inventory.Inventory, inventor
if _, err := inv.AddNode(ctx, report.Node); err != nil { if _, err := inv.AddNode(ctx, report.Node); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, report); err != nil { if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, report); err != nil {
t.Fatal(err) t.Fatal(err)
} }
doing, said, err := inv.DoingOf(ctx, report.Node) doing, said, err := inv.DoingOf(ctx, report.Node)
@@ -103,7 +103,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T)
if err := inv.RecordSent(ctx, node.ID, digest); err != nil { if err := inv.RecordSent(ctx, node.ID, digest); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{ if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{
Node: "anchor", Applied: []string{"a", "b"}, Declared: digest, Carried: []int{5432}, Node: "anchor", Applied: []string{"a", "b"}, Declared: digest, Carried: []int{5432},
}); err != nil { }); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -126,7 +126,7 @@ func TestABareAliveDoesNotWipeTheDeclarationThatSaysANodeIsCurrent(t *testing.T)
} }
// Now the node says only that it is there, as it does every minute. // Now the node says only that it is there, as it does every minute.
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil { if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if !currentOf("anchor") { if !currentOf("anchor") {
@@ -162,7 +162,7 @@ func TestAFailureDoesNotBecomeTheAccountOfWhatTheMachineHolds(t *testing.T) {
if err := inv.RecordOwned(ctx, node.ID, []string{"one", "two", "three"}); err != nil { if err := inv.RecordOwned(ctx, node.ID, []string{"one", "two", "three"}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{ if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{
Node: "workstation", Applied: []string{"one"}, Failed: map[string]string{"two": "no"}, Node: "workstation", Applied: []string{"one"}, Failed: map[string]string{"two": "no"},
}); err != nil { }); err != nil {
t.Fatal(err) t.Fatal(err)
@@ -200,13 +200,13 @@ func TestWhatAnAdoptedNodeHoldsIsKeptAndAnAliveWordDoesNotWipeIt(t *testing.T) {
} }
check("after the report") check("after the report")
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil { if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor"}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
check("after an alive word") check("after an alive word")
// A reconcile report carrying only adoption is recorded, though it applied nothing. // A reconcile report carrying only adoption is recorded, though it applied nothing.
if err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor", if _, err := (link.Enrolment{Inventory: inv}).Heard(ctx, link.Report{Node: "anchor",
Firewall: "ufw"}); err != nil { Firewall: "ufw"}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
+37
View File
@@ -154,10 +154,47 @@ func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Con
m.seq = meta.Sequence.Stream m.seq = meta.Sequence.Stream
m.delivered = meta.NumDelivered m.delivered = meta.NumDelivered
} }
// **Work that outlives the acknowledgement window says so while it runs.**
//
// The bus waits a fixed time to be told a message was taken, and then hands it to whoever
// consumes next — which is right for a consumer that died and wrong for one that is busy.
// Acting on a merge builds every module the merge changed: minutes of work against a
// thirty-second window. So the same merge was handed over again while the first build was
// still running, and again after that — on 2026-09-28 one merge ran the mesh's whole
// catalogue five times over and exhausted a public registry's pull limit.
//
// Here rather than in each handler, because the window belongs to the transport and every
// handler would otherwise have to remember it. It changes nothing about a handler that
// dies: a message is kept alive only while this goroutine is, so a controller that stops
// stops saying so, and the bus redelivers exactly as it should.
working := make(chan struct{})
defer close(working)
go stillWorking(msg, working)
} }
act(ctx, m) act(ctx, m)
} }
// heartbeatWhileWorking is how often a handler still running tells the bus so — comfortably inside
// the shortest acknowledgement window the mesh gives any of its consumers.
const heartbeatWhileWorking = 10 * time.Second
// stillWorking keeps one message alive until the work on it returns.
//
// An error is not worth reporting: what the bus does when it is not told is redeliver, which is
// exactly what happens if this fails, and the handler's own outcome is the thing worth logging.
func stillWorking(msg *nats.Msg, done <-chan struct{}) {
tick := time.NewTicker(heartbeatWhileWorking)
defer tick.Stop()
for {
select {
case <-done:
return
case <-tick.C:
_ = msg.InProgress()
}
}
}
// kindOfSubject is how this transport's addressing becomes what the mesh calls a message. // kindOfSubject is how this transport's addressing becomes what the mesh calls a message.
// //
// By subject, which is the only thing the server enforces: a body claiming to be a report does not // By subject, which is the only thing the server enforces: a body claiming to be a report does not
+167 -5
View File
@@ -86,14 +86,15 @@ type counted struct {
heard []Report heard []Report
} }
func (c *counted) Heard(_ context.Context, r Report) error { func (c *counted) Heard(_ context.Context, r Report) (bool, error) {
c.mu.Lock() c.mu.Lock()
defer c.mu.Unlock() defer c.mu.Unlock()
if c.err != nil { if c.err != nil {
return c.err return false, c.err
} }
c.heard = append(c.heard, r) c.heard = append(c.heard, r)
return nil // News, so what the mesh states about a report is exercised wherever a report is.
return true, nil
} }
func (c *counted) refusing(err error) { func (c *counted) refusing(err error) {
@@ -207,11 +208,11 @@ type sentAndHeardSafely struct {
heard []Report heard []Report
} }
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) error { func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) (bool, error) {
s.mu.Lock() s.mu.Lock()
defer s.mu.Unlock() defer s.mu.Unlock()
s.heard = append(s.heard, r) s.heard = append(s.heard, r)
return nil return true, nil
} }
func (s *sentAndHeardSafely) Outstanding(context.Context, string) (string, error) { func (s *sentAndHeardSafely) Outstanding(context.Context, string) (string, error) {
@@ -395,3 +396,164 @@ func TestEverySubjectTheControllerFollowsDecodesToAKind(t *testing.T) {
t.Error("a subject nobody follows decoded to a kind") t.Error("a subject nobody follows decoded to a kind")
} }
} }
// **A handler slower than the acknowledgement window is not handed its message again.**
//
// The bus waits a fixed time to be told a message was taken and then redelivers, which is right for
// a consumer that died and wrong for one that is busy. Acting on a merge builds modules — minutes
// against a thirty-second window — and the same merge was handed over five times while the first
// build was still running (2026-09-28). Here the window is two seconds and the work takes six.
func TestNatsWorkSlowerThanTheWindowIsNotHandedOverAgain(t *testing.T) {
js := aBus(t)
// The controller's own consumer, with a window short enough to outlive in a test.
if err := js.EnsureConsumer(broker.Consumer{
Name: broker.ControllerName, Stream: "CONTROL", Push: true, AckWaitSeconds: 2,
Why: "a window short enough to outlive in a test",
}); err != nil {
t.Fatal(err)
}
var mu sync.Mutex
handled := 0
slow := make(chan struct{})
s, stop := servingOn(t, js, nil)
defer stop()
s.listener = slowly{func() {
mu.Lock()
handled++
first := handled == 1
mu.Unlock()
if first {
time.Sleep(6 * time.Second)
close(slow)
}
}}
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
select {
case <-slow:
case <-time.After(30 * time.Second):
t.Fatal("the slow work never finished")
}
// A moment for a redelivery to arrive, if the bus were going to send one.
time.Sleep(3 * time.Second)
mu.Lock()
defer mu.Unlock()
if handled != 1 {
t.Fatalf("one report was handled %d times, so slow work is run again while it is running", handled)
}
}
// slowly is a listener that runs whatever it was given.
type slowly struct{ work func() }
func (s slowly) Heard(context.Context, Report) (bool, error) { s.work(); return true, nil }
// **The mesh says what it applied** (novox/hq ADR 0134), under the seat the control plane holds — and
// says nothing when a report is the same state said again, which is what a machine reconciling every
// minute sends.
func TestNatsTheMeshSaysWhatAMachineApplied(t *testing.T) {
js := aBus(t)
heard := make(chan *nats.Msg, 4)
sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) {
heard <- m
})
if err != nil {
t.Fatal(err)
}
defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection
_, stop := servingOn(t, js, &counted{})
defer stop()
// A report that changed something: the store says it was news.
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store", "broker"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
select {
case m := <-heard:
if m.Subject != SeatEventSubject(MeshControllerSeat, KeyApplied) {
t.Fatalf("the mesh stated %q", m.Subject)
}
var said Applied
if err := json.Unmarshal(m.Data, &said); err != nil {
t.Fatal(err)
}
if said.Node != "anchor" || said.Declared != "d1" || said.Resources != 2 {
t.Fatalf("it said %+v", said)
}
case <-time.After(10 * time.Second):
t.Fatal("the mesh said nothing about a machine that now runs something else")
}
// A refusal is its own fact, with the reason in it rather than only in a log.
refusal, err := json.Marshal(Report{Node: "anchor", Declared: "d2",
Failed: map[string]string{"gitea.server": "no such image"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), refusal); err != nil {
t.Fatal(err)
}
select {
case m := <-heard:
if m.Subject != SeatEventSubject(MeshControllerSeat, KeyRefused) {
t.Fatalf("a refusal was stated as %q", m.Subject)
}
var said Refused
if err := json.Unmarshal(m.Data, &said); err != nil {
t.Fatal(err)
}
if said.Failed["gitea.server"] == "" {
t.Fatalf("the refusal does not say which resource or why: %+v", said)
}
case <-time.After(10 * time.Second):
t.Fatal("the mesh said nothing about a machine that refused what it was sent")
}
}
// And a report that is not news is not a fact. A machine reconciles every minute; a fact per report
// would be a fact per minute per machine, which is a stream nobody reads.
func TestNatsAReportThatIsNotNewsIsNotStated(t *testing.T) {
js := aBus(t)
heard := make(chan *nats.Msg, 4)
sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) {
heard <- m
})
if err != nil {
t.Fatal(err)
}
defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection
// A store that records the report and says it was nothing new — which is what the mesh's own
// store says about a machine repeating itself.
_, stop := servingOn(t, js, sameAgain{})
defer stop()
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
select {
case m := <-heard:
t.Fatalf("the mesh stated %q about a machine that changed nothing", m.Subject)
case <-time.After(3 * time.Second):
}
}
// sameAgain records a report and says it was the same state said again.
type sameAgain struct{}
func (sameAgain) Heard(context.Context, Report) (bool, error) { return false, nil }
+4 -4
View File
@@ -60,7 +60,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) {
rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()} rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()}
rekey.Proof = ed25519.Sign(private, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel())) rekey.Proof = ed25519.Sign(private, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel()))
if err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}); err != nil { if _, err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}); err != nil {
t.Fatal(err) t.Fatal(err)
} }
placed, err := e.Inventory.Overlays(ctx) placed, err := e.Inventory.Overlays(ctx)
@@ -77,7 +77,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) {
_ = hub _ = hub
// Replayed, it is stale: the previous key it names is no longer the node's. // Replayed, it is stale: the previous key it names is no longer the node's.
err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}) _, err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey})
if err == nil || !strings.Contains(err.Error(), "previous overlay key") { if err == nil || !strings.Contains(err.Error(), "previous overlay key") {
t.Fatalf("a replayed rekey was accepted: %v", err) t.Fatalf("a replayed rekey was accepted: %v", err)
} }
@@ -93,7 +93,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) {
rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()} rekey := &link.Rekey{Previous: ownKey, OverlayKey: tunnelKey, Tunnel: theTunnel()}
rekey.Proof = ed25519.Sign(stranger, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel())) rekey.Proof = ed25519.Sign(stranger, link.RekeyProof("anchor", ownKey, tunnelKey, theTunnel()))
err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey}) _, err = e.Heard(ctx, link.Report{Node: "anchor", Rekey: rekey})
if err == nil || !strings.Contains(err.Error(), "not signed by anchor's identity key") { if err == nil || !strings.Contains(err.Error(), "not signed by anchor's identity key") {
t.Fatalf("a rekey signed by a stranger was accepted: %v", err) t.Fatalf("a rekey signed by a stranger was accepted: %v", err)
} }
@@ -111,7 +111,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) {
other := theTunnel() other := theTunnel()
other.Port = 51820 other.Port = 51820
moved.Proof = ed25519.Sign(mustPrivate(t, e, "anchor"), link.RekeyProof("anchor", ownKey, tunnelKey, other)) moved.Proof = ed25519.Sign(mustPrivate(t, e, "anchor"), link.RekeyProof("anchor", ownKey, tunnelKey, other))
if err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: moved}); err == nil { if _, err := e.Heard(ctx, link.Report{Node: "anchor", Rekey: moved}); err == nil {
t.Fatal("a proof over another tunnel was accepted") t.Fatal("a proof over another tunnel was accepted")
} }
} }
+2 -2
View File
@@ -9,12 +9,12 @@ import (
type heardWith struct{ err error } type heardWith struct{ err error }
func (h heardWith) Heard(context.Context, Report) error { return h.err } func (h heardWith) Heard(context.Context, Report) (bool, error) { return h.err == nil, h.err }
// switchable answers with whatever it is set to — the store away, then back. // switchable answers with whatever it is set to — the store away, then back.
type switchable struct{ err error } type switchable struct{ err error }
func (h *switchable) Heard(context.Context, Report) error { return h.err } func (h *switchable) Heard(context.Context, Report) (bool, error) { return h.err == nil, h.err }
func aReport(node, declared string) Report { func aReport(node, declared string) Report {
return Report{Node: node, Declared: declared, Applied: []string{"store"}} return Report{Node: node, Declared: declared, Applied: []string{"store"}}
+62 -4
View File
@@ -29,7 +29,12 @@ type Enroller interface {
// Listener is what the controller does with a report. Separate from Enroller so the two can be // Listener is what the controller does with a report. Separate from Enroller so the two can be
// given independently, and so a server that only sends declarations needs neither. // given independently, and so a server that only sends declarations needs neither.
type Listener interface { type Listener interface {
Heard(ctx context.Context, report Report) error // Heard records what a node said, and says whether it was **news** — a machine that now runs
// something else, or refuses something it did not refuse before. A machine reconciles
// continuously and reports each time, so what is news is the store's answer rather than the
// bus's: only this side has the previous report to compare with. What the mesh states about it
// is the server's (novox/hq ADR 0134).
Heard(ctx context.Context, report Report) (news bool, err error)
} }
// Recorder keeps what builders say. // Recorder keeps what builders say.
@@ -257,7 +262,7 @@ func (s *Server) heartbeat(m Control) {
return return
} }
if s.listener != nil { if s.listener != nil {
if err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil { if _, err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil {
s.log.Printf("could not record that %s is here: %v", alive.Node, err) s.log.Printf("could not record that %s is here: %v", alive.Node, err)
} }
} }
@@ -291,7 +296,7 @@ func (s *Server) reported(ctx context.Context, m Control) {
return return
} }
err := s.listener.Heard(context.Background(), report) news, err := s.listener.Heard(context.Background(), report)
switch s.decide(ctx, m, what, declaredIn, outstanding, err) { switch s.decide(ctx, m, what, declaredIn, outstanding, err) {
case Hold: case Hold:
// Held, not settled, while the store cannot take it: the node reports an apply once, // Held, not settled, while the store cannot take it: the node reports an apply once,
@@ -307,6 +312,13 @@ func (s *Server) reported(ctx context.Context, m Control) {
// node whose recovery copy is silently older than it looks. // node whose recovery copy is silently older than it looks.
s.log.Printf("could not record %s's report: %v", report.Node, err) s.log.Printf("could not record %s's report: %v", report.Node, err)
} }
// **And the mesh says what it did** (novox/hq ADR 0134). Only when the report was news: a
// machine reports every convergence, and a fact per report would be a fact per minute per
// machine saying nothing. Stated after it is recorded, so nothing is announced that the
// mesh does not hold.
if err == nil && news {
s.saysWhatItDid(ctx, report)
}
} }
switch { switch {
@@ -460,7 +472,19 @@ func (s *Server) catchingUp(ctx context.Context, m Control) {
sent := 0 sent := 0
for _, a := range announcements { for _, a := range announcements {
a.Replay = true a.Replay = true
if err := EmitEvent(ctx, s.bus, KeyModuleBuilt, "control-plane", "", a); err != nil { // Under the control plane's own seat (novox/hq ADR 0134). It used to be published as a
// module's event from a module called "control-plane", which does not exist — so the
// controller's own account refused it, every catalogue that asked what it missed was
// answered with nothing, and its graph kept the gap (found 2026-09-28).
body, err := json.Marshal(a)
if err != nil {
// A body that cannot be written is this program's fault, not the bus's, and publishing
// an empty one would put a fact on the mesh that says nothing.
s.log.Printf("cannot re-announce %s at %s: %v", a.Module, short(a.Commit), err)
_ = m.Took()
return
}
if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, KeyBuiltBefore, body); err != nil {
// Said and abandoned rather than retried: the catalogue asks again every time it // Said and abandoned rather than retried: the catalogue asks again every time it
// starts, and half a graph delivered twice is no better than half delivered once. // starts, and half a graph delivered twice is no better than half delivered once.
s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v", s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v",
@@ -551,3 +575,37 @@ func (s *Server) sourceMoved(ctx context.Context, m Control) {
} }
_ = m.Took() _ = m.Took()
} }
// saysWhatItDid states what a machine now runs, or what it would not take, as a fact on the bus
// (novox/hq ADR 0134).
//
// **The control plane speaks, as the holder of its seat.** A node's report is control traffic only
// this process may read, so the chain from a merge to a machine went dark exactly where it touched
// one: nothing said which version a machine runs, or that it refused to. The facts are second-hand
// on purpose — one emitter, one ordering — and a machine that cannot reach the bus produces none, so
// absence is not health.
//
// A failure to state a fact is logged and nothing else: the report is recorded, which is the part
// that must not be lost, and the next change says the same thing again.
func (s *Server) saysWhatItDid(ctx context.Context, report Report) {
if s.bus == nil {
return
}
event, body := KeyApplied, any(Applied{
Node: report.Node, Declared: report.Declared, Resources: len(report.Applied),
})
if report.Refused != "" || len(report.Failed) > 0 {
event, body = KeyRefused, Refused{
Node: report.Node, Declared: report.Declared,
Refused: report.Refused, Failed: report.Failed,
}
}
raw, err := json.Marshal(body)
if err != nil {
s.log.Printf("could not say what %s did: %v", report.Node, err)
return
}
if err := s.bus.PublishSeatEvent(ctx, MeshControllerSeat, event, raw); err != nil {
s.log.Printf("could not say that %s %s: %v", report.Node, event, err)
}
}
+3 -3
View File
@@ -23,12 +23,12 @@ type sentAndHeard struct {
err error err error
} }
func (s *sentAndHeard) Heard(_ context.Context, r Report) error { func (s *sentAndHeard) Heard(_ context.Context, r Report) (bool, error) {
if s.err != nil { if s.err != nil {
return s.err return false, s.err
} }
s.heard = append(s.heard, r) s.heard = append(s.heard, r)
return nil return true, nil
} }
func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil } func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil }
+1 -24
View File
@@ -27,6 +27,7 @@
"bus": "/var/lib/mesh/mesh-controller/bus" "bus": "/var/lib/mesh/mesh-controller/bus"
}, },
"secrets-owner": "65534:65534", "secrets-owner": "65534:65534",
"prepares": true,
"resources": [ "resources": [
{ {
"id": "mesh-state", "id": "mesh-state",
@@ -34,30 +35,6 @@
"path": "/var/lib/mesh/mesh-controller", "path": "/var/lib/mesh/mesh-controller",
"mode": "0700" "mode": "0700"
}, },
{
"id": "migrate",
"type": "container",
"name": "mesh-controller-migrate",
"run-once": true,
"network": "host",
"args": [
"migrate"
],
"env": {
"MESH_STORE_INVENTORY_FILE": "/run/secrets/inventory",
"MESH_STORE_IDENTITY_FILE": "/run/secrets/identity",
"MESH_STORE_LICENCES_FILE": "/run/secrets/licences",
"MESH_STORE_INVENTORY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_IDENTITY_PORT": "${seat:mesh-store:5432}",
"MESH_STORE_LICENCES_PORT": "${seat:mesh-store:5432}"
},
"volumes": [
"/var/lib/mesh/mesh-controller/inventory:/run/secrets/inventory:ro",
"/var/lib/mesh/mesh-controller/identity:/run/secrets/identity:ro",
"/var/lib/mesh/mesh-controller/licences:/run/secrets/licences:ro"
],
"artifact": "server"
},
{ {
"id": "server", "id": "server",
"type": "container", "type": "container",