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
33 changed files with 781 additions and 100 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.
+3
View File
@@ -148,6 +148,9 @@ func usage() {
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 -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
} }
+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)
+10 -2
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 {
@@ -1685,7 +1690,10 @@ func prepared(from map[string]any) map[string]any {
for k, v := range from { for k, v := range from {
step[k] = v step[k] = v
} }
step["id"] = fmt.Sprint(from["id"]) + ".prepare" // **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["name"] = fmt.Sprint(from["name"]) + "-prepare"
step["run-once"] = true step["run-once"] = true
step["args"] = []any{PreparationArgument} step["args"] = []any{PreparationArgument}
+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)
} }
+9 -3
View File
@@ -3,6 +3,7 @@ package catalogue
import ( import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"strings"
"testing" "testing"
) )
@@ -57,13 +58,18 @@ func TestThePreparationRunsTheModulesOwnCodeAndComesRightBeforeIt(t *testing.T)
ids := idsOf(out) ids := idsOf(out)
at := -1 at := -1
for i, id := range ids { for i, id := range ids {
if id == "gitea.runtime.prepare" { if id == "gitea.runtime-prepare" {
at = i at = i
} }
} }
if at < 0 { if at < 0 {
t.Fatalf("nothing prepares this module's state: %v", ids) 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" { if ids[at+1] != "gitea.runtime" {
t.Fatalf("the preparation is not immediately before the module's own code: %v", ids) t.Fatalf("the preparation is not immediately before the module's own code: %v", ids)
} }
@@ -79,7 +85,7 @@ func TestThePreparationRunsTheModulesOwnCodeAndComesRightBeforeIt(t *testing.T)
func TestThePreparationIsGivenWhatTheModuleIsGiven(t *testing.T) { func TestThePreparationIsGivenWhatTheModuleIsGiven(t *testing.T) {
out := declaredFor(t, aPreparingModule()) out := declaredFor(t, aPreparingModule())
declared := byID(out) declared := byID(out)
step, workload := declared["gitea.runtime.prepare"], declared["gitea.runtime"] step, workload := declared["gitea.runtime-prepare"], declared["gitea.runtime"]
if step == nil || workload == nil { if step == nil || workload == nil {
t.Fatalf("expected both, got %v", idsOf(out)) t.Fatalf("expected both, got %v", idsOf(out))
} }
@@ -108,7 +114,7 @@ func TestAModuleThatPreparesNothingGetsNoStep(t *testing.T) {
m := aPreparingModule() m := aPreparingModule()
m.Prepares = false m.Prepares = false
for _, id := range idsOf(declaredFor(t, m)) { for _, id := range idsOf(declaredFor(t, m)) {
if id == "gitea.runtime.prepare" { if id == "gitea.runtime-prepare" {
t.Fatal("a module that prepares nothing was given a preparation") t.Fatal("a module that prepares nothing was given a preparation")
} }
} }
+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)
} }
+111 -6
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) {
@@ -451,4 +452,108 @@ func TestNatsWorkSlowerThanTheWindowIsNotHandedOverAgain(t *testing.T) {
// slowly is a listener that runs whatever it was given. // slowly is a listener that runs whatever it was given.
type slowly struct{ work func() } type slowly struct{ work func() }
func (s slowly) Heard(context.Context, Report) error { s.work(); return nil } 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 }