Compare commits
19
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6215ff0760 | ||
|
|
54812306be | ||
|
|
ce9e20fbbc | ||
|
|
878690697e | ||
|
|
ad97297576 | ||
|
|
683b1ed693 | ||
|
|
04f9f378b0 | ||
|
|
1c3f44a526 | ||
|
|
89e152dfe2 | ||
|
|
1ebad3786c | ||
|
|
f2f526a60a | ||
|
|
4b4c7e0e0d | ||
|
|
cec792ce9d | ||
|
|
338d033632 | ||
|
|
1be926cec4 | ||
|
|
2134768dfe | ||
|
|
ef825688ee | ||
|
|
77a14360df | ||
|
|
9be2fb4750 |
@@ -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 {
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
@@ -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" }
|
||||||
} }
|
} }
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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}
|
||||||
|
|||||||
@@ -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]
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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))
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
@@ -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
@@ -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.
|
||||||
|
|||||||
@@ -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
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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"
|
||||||
|
|||||||
@@ -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)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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 }
|
||||||
|
|||||||
@@ -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")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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 }
|
||||||
|
|||||||
Reference in New Issue
Block a user