Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
fe5988c536 | ||
|
|
ed5d467d90 | ||
|
|
228d0226dd | ||
|
|
6215ff0760 | ||
|
|
54812306be | ||
|
|
ce9e20fbbc | ||
|
|
878690697e | ||
|
|
ad97297576 | ||
|
|
683b1ed693 | ||
|
|
04f9f378b0 | ||
|
|
1c3f44a526 | ||
|
|
89e152dfe2 | ||
|
|
1ebad3786c | ||
|
|
f2f526a60a |
@@ -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 {
|
||||
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),
|
||||
Firewall: "ufw", Held: held, Reachable: reachable,
|
||||
}); err != nil {
|
||||
|
||||
@@ -309,7 +309,7 @@ func converge(ctx context.Context, open *stores, node string, yes bool, digest s
|
||||
return "", err
|
||||
}
|
||||
derived := derivedFilter{rules: rules, foundation: with.Foundation, mesh: with.Mesh,
|
||||
outward: plan.PublicDomain != ""}
|
||||
outward: plan.PublicDomain != "", outwardLinks: with.OutwardLinks}
|
||||
preview, saw := previewOf(node, reported, derived, plan, taken, filter, runs[filter])
|
||||
preview += "\n\n preview " + saw
|
||||
if !yes {
|
||||
@@ -414,6 +414,18 @@ func previewOf(node string, reported inventory.Adoption, derived derivedFilter,
|
||||
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 " +
|
||||
"declares it\n")
|
||||
// Which links the filter constrains, said rather than left to the sentence above (novox/hq ADR
|
||||
// 0140). Everything arriving anywhere else is this machine's own guest and keeps working — which
|
||||
// is what a reader most wants to know, because the previous shape of this filter cut a machine's
|
||||
// guests off at the flip without saying so, and that is how this was found.
|
||||
if len(derived.outwardLinks) > 0 {
|
||||
b.WriteString(fmt.Sprintf(" it filters what arrives on: %s, and on the private network "+
|
||||
"— everything its own guests send keeps working\n",
|
||||
strings.Join(derived.outwardLinks, ", ")))
|
||||
} else {
|
||||
b.WriteString(" it has reported no link facing outside, so no filter can be composed " +
|
||||
"for it — the flip is refused until it reports one\n")
|
||||
}
|
||||
|
||||
isTaken := map[string]bool{}
|
||||
for _, m := range taken {
|
||||
@@ -473,6 +485,10 @@ type derivedFilter struct {
|
||||
// mesh is every address on the private network; outward says the machine faces outside.
|
||||
mesh []string
|
||||
outward bool
|
||||
// outwardLinks is the links this machine reported as facing outside it (novox/hq ADR 0140).
|
||||
// The filter constrains what arrives on them; everything arriving elsewhere is this machine's
|
||||
// own guest and is not filtered.
|
||||
outwardLinks []string
|
||||
}
|
||||
|
||||
// closesOutside is what a narrowing from everywhere to the private network is called: it closes.
|
||||
@@ -498,6 +514,12 @@ func (d derivedFilter) fate(r inventory.Reach) string {
|
||||
return "stays open — the mesh's own, from anywhere"
|
||||
}
|
||||
}
|
||||
// This machine's own guests ask it for an address and for names, and those two arrive here
|
||||
// (novox/hq ADR 0140). Admitted by the link they arrive on, so a listener bound anywhere but an
|
||||
// outward link keeps answering them.
|
||||
if (r.Protocol == "udp" && (r.Port == 53 || r.Port == 67)) || (r.Protocol == "tcp" && r.Port == 53) {
|
||||
return "stays open — this machine's own guests asking it for an address and for names"
|
||||
}
|
||||
for _, rule := range d.rules {
|
||||
if rule.Port != r.Port || rule.Protocol != r.Protocol {
|
||||
continue
|
||||
|
||||
@@ -148,6 +148,9 @@ func usage() {
|
||||
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> --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 --new <name> create the record and issue for it
|
||||
token issue ... --adopted ...for a machine in use, which joins adopted
|
||||
|
||||
@@ -66,6 +66,18 @@ func nodeCommand(ctx context.Context, args []string) error {
|
||||
// because the damage is already done by the time it prints.
|
||||
return publicDomain(ctx, inv, args[1:])
|
||||
|
||||
case "networks":
|
||||
// Removed by novox/hq ADR 0140, which superseded the record that added it. The filter no
|
||||
// longer names any network: it constrains what arrives from outside the machine and says
|
||||
// nothing about what did not, so there is no list to keep. Answered rather than met with
|
||||
// "unknown command", because this was the documented way to stop a flip cutting a machine's
|
||||
// containers off and somebody will reasonably still type it.
|
||||
return errors.New("`node networks` is gone (novox/hq ADR 0140). The filter constrains what " +
|
||||
"arrives from outside this machine and says nothing about traffic that did not, so no " +
|
||||
"network is named anywhere and nothing needs to be said to keep a machine's own " +
|
||||
"containers reaching outward. The machine reports which of its links face outside; see " +
|
||||
"`node show <name>`")
|
||||
|
||||
case "account":
|
||||
// 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;
|
||||
|
||||
@@ -640,13 +640,19 @@ func renderingFor(ctx context.Context, open *stores, node string,
|
||||
if err != nil {
|
||||
return catalogue.Rendering{}, inventory.Node{}, err
|
||||
}
|
||||
// Which of this machine's links face outside, which is what the derived filter is written
|
||||
// around (novox/hq ADR 0140). Reported by the machine, never set.
|
||||
outwardLinks, err := inv.OutwardLinksOf(ctx, node)
|
||||
if err != nil {
|
||||
return catalogue.Rendering{}, inventory.Node{}, err
|
||||
}
|
||||
return catalogue.Rendering{
|
||||
BusMembership: memberships[node],
|
||||
Settings: settings, Generators: gens, Grants: grants, Needed: needed, Ports: ports,
|
||||
Certificate: certificate, Authority: authority, Mesh: private, Names: names,
|
||||
Machines: machines,
|
||||
Suffix: overlay.Suffix(), MeshRange: meshRange, Accounts: accounts, Foundation: foundation,
|
||||
Kept: kept, Adopted: record.Adopted,
|
||||
Suffix: overlay.Suffix(), MeshRange: meshRange, TunnelInterface: overlay.Interface, Accounts: accounts, Foundation: foundation,
|
||||
Kept: kept, Adopted: record.Adopted, OutwardLinks: outwardLinks,
|
||||
Given: given, Taken: taken, Seats: seats, ArtifactStore: artifactStore, Built: built,
|
||||
BusUsers: busUsers,
|
||||
}, record, nil
|
||||
|
||||
@@ -162,7 +162,13 @@ func declare(ctx context.Context, args []string) error {
|
||||
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 {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -181,6 +181,14 @@ func PermissionsFor(p Principal) (Permissions, error) {
|
||||
for _, seat := range meshSeatsTheControllerUses {
|
||||
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
|
||||
// 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
|
||||
|
||||
@@ -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.>`).
|
||||
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
|
||||
// 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
|
||||
users = [
|
||||
{ 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"] }
|
||||
allow_responses: { max: 1, ttl: "1m" }
|
||||
} }
|
||||
|
||||
@@ -137,7 +137,7 @@ func (r Registry) has(ctx context.Context, url string, accept ...string) (bool,
|
||||
return false, err
|
||||
}
|
||||
for _, media := range accept {
|
||||
request.Header.Set("Accept", media)
|
||||
request.Header.Add("Accept", media)
|
||||
}
|
||||
response, err := r.client().Do(request)
|
||||
if err != nil {
|
||||
|
||||
@@ -59,7 +59,11 @@ func anchorRendering(adopted bool) Rendering {
|
||||
Values: map[string]any{ExposeSetting: map[string]any{"5000": FromEverywhere}}}}},
|
||||
Mesh: []string{"10.42.0.1"},
|
||||
Foundation: []int{5671},
|
||||
Adopted: adopted,
|
||||
// What the machine reported faces outside, which every rule in the filter is written
|
||||
// around (novox/hq ADR 0140).
|
||||
OutwardLinks: []string{"eth0"},
|
||||
TunnelInterface: "mesh0",
|
||||
Adopted: adopted,
|
||||
// Genesis takes the foundation's modules.
|
||||
Taken: map[string]bool{"postgres": true, "lavinmq": true},
|
||||
}
|
||||
@@ -575,11 +579,13 @@ func TestAGivenMachineSideReachesTheFilterTheOpeningAndTheConsumer(t *testing.T)
|
||||
}
|
||||
r := Resolution{Node: "anchor", Modules: []Manifest{forge}}
|
||||
with := Rendering{
|
||||
Ports: map[string]map[int]int{"forge": portsAsThePlanWould(forge, given)},
|
||||
Given: map[string]map[int]int{"forge": given},
|
||||
Mesh: []string{"10.77.0.1"},
|
||||
Adopted: true,
|
||||
Taken: map[string]bool{"forge": true},
|
||||
Ports: map[string]map[int]int{"forge": portsAsThePlanWould(forge, given)},
|
||||
Given: map[string]map[int]int{"forge": given},
|
||||
Mesh: []string{"10.77.0.1"},
|
||||
Adopted: true,
|
||||
OutwardLinks: []string{"eth0"},
|
||||
TunnelInterface: "mesh0",
|
||||
Taken: map[string]bool{"forge": true},
|
||||
}
|
||||
|
||||
// What the runtime is handed: the machine's own port on the outside, the container's within.
|
||||
@@ -660,9 +666,11 @@ func TestALongFormPortIsOpenedWhereTheManifestPublishesIt(t *testing.T) {
|
||||
forge := aForge()
|
||||
r := Resolution{Node: "anchor", Modules: []Manifest{forge}}
|
||||
composed, err := r.Compose(Rendering{
|
||||
Ports: map[string]map[int]int{"forge": portsAsThePlanWould(forge, nil)},
|
||||
Mesh: []string{"10.77.0.1"},
|
||||
Adopted: true,
|
||||
Ports: map[string]map[int]int{"forge": portsAsThePlanWould(forge, nil)},
|
||||
Mesh: []string{"10.77.0.1"},
|
||||
Adopted: true,
|
||||
OutwardLinks: []string{"eth0"},
|
||||
TunnelInterface: "mesh0",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -153,6 +153,26 @@ func (b *Build) problems(module string) []string {
|
||||
"%s: %q is a bundle and says no language, so nothing can choose a compiler "+
|
||||
"for it", module, a.Name))
|
||||
}
|
||||
// **A system, for a language that compiles to a binary** (novox/hq ADR 0142). A binary
|
||||
// is pinned to one operating system at link time so a host refuses to touch a machine
|
||||
// it was not built for (novox/hq ADR 0005); an artifact that says nothing would be
|
||||
// compiled for whatever the build machine happened to be, which reads as portable and
|
||||
// is not.
|
||||
if compiled := compilesToABinary(a.Language); compiled && strings.TrimSpace(a.System) == "" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s: %q is compiled to a binary and says no system, so it would be built for "+
|
||||
"whatever the build machine happens to be. Declare one artifact per "+
|
||||
"system: %s", module, a.Name, spokenSystems()))
|
||||
} else if !compiled && strings.TrimSpace(a.System) != "" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s: %q names the system %q and is written in %q, which compiles to code that "+
|
||||
"runs anywhere — a system that decides nothing reads as though it did",
|
||||
module, a.Name, a.System, a.Language))
|
||||
} else if compiled && !knownSystem(a.System) {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
"%s: %q is built for %q, and a system is %s",
|
||||
module, a.Name, a.System, spokenSystems()))
|
||||
}
|
||||
} else {
|
||||
if a.From == "" {
|
||||
problems = append(problems, fmt.Sprintf(
|
||||
@@ -193,3 +213,42 @@ func oneOrOther(n int) string {
|
||||
}
|
||||
return "them"
|
||||
}
|
||||
|
||||
// Systems the mesh builds binaries for, which is the set a host may be pinned to (novox/hq ADR 0005).
|
||||
//
|
||||
// **A closed list, and the host's own, not the compiler's.** These are not the values a Go toolchain
|
||||
// would call an operating system — the difference between two of them is a C library, not a kernel.
|
||||
// They are what a machine reports itself to be and what a host is linked to refuse, so the list that
|
||||
// matters is the one the host understands.
|
||||
var systems = []string{"alpine", "android", "arch"}
|
||||
|
||||
// knownSystem is whether the mesh builds for it.
|
||||
func knownSystem(system string) bool {
|
||||
want := strings.ToLower(strings.TrimSpace(system))
|
||||
for _, s := range systems {
|
||||
if s == want {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// spokenSystems is the list as a refusal says it, so a reader is one edit from right.
|
||||
func spokenSystems() string {
|
||||
return strings.Join(systems, ", ")
|
||||
}
|
||||
|
||||
// compilesToABinary is whether this language's bundle is a binary for one operating system rather
|
||||
// than code that runs wherever its interpreter does.
|
||||
//
|
||||
// **Asked of the language, not of the artifact.** A module says what it is written in; what that
|
||||
// implies is the mesh's to know, exactly as the compiler is (novox/hq ADR 0142). Asking the artifact
|
||||
// would let two artifacts in one language disagree about whether they are portable.
|
||||
func compilesToABinary(language string) bool {
|
||||
switch strings.ToLower(strings.TrimSpace(language)) {
|
||||
case "go":
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
package catalogue
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// bundleFor is a manifest whose one artifact is a bundle in the given language and system.
|
||||
func bundleFor(language, system string) Manifest {
|
||||
return Manifest{Module: "a-component", Build: &Build{Artifacts: []Artifact{
|
||||
{Name: "binary", Kind: ArtifactBundle, Language: language, System: system},
|
||||
}}}
|
||||
}
|
||||
|
||||
func problemsOf(t *testing.T, m Manifest) string {
|
||||
t.Helper()
|
||||
return strings.Join(m.Build.problems(m.Module), "\n")
|
||||
}
|
||||
|
||||
// **A language that compiles to a binary must say which system.**
|
||||
//
|
||||
// A binary is pinned to one operating system at link time, so a host refuses to touch a machine it
|
||||
// was not built for. An artifact that says nothing would be compiled for whatever the build machine
|
||||
// happened to be — which reads as portable and is not, and is the fault this check exists for.
|
||||
func TestABinaryMustSayWhichSystemItIsFor(t *testing.T) {
|
||||
got := problemsOf(t, bundleFor("go", ""))
|
||||
if !strings.Contains(got, "says no system") {
|
||||
t.Fatalf("a compiled bundle with no system was accepted:\n%s", got)
|
||||
}
|
||||
// And the refusal names what it could have said, so a reader is one edit from right.
|
||||
for _, system := range []string{"alpine", "android", "arch"} {
|
||||
if !strings.Contains(got, system) {
|
||||
t.Fatalf("the refusal does not name %q as a choice:\n%s", system, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestABinaryThatNamesASystemIsAccepted(t *testing.T) {
|
||||
if got := problemsOf(t, bundleFor("go", "arch")); got != "" {
|
||||
t.Fatalf("a compiled bundle naming a system was refused:\n%s", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A system the mesh does not build for is refused where it is written. These are the host's own
|
||||
// names, not a compiler's: the difference between two of them is a C library rather than a kernel,
|
||||
// so a value that looks like an operating system to a toolchain is still wrong here.
|
||||
func TestASystemTheMeshDoesNotBuildForIsRefused(t *testing.T) {
|
||||
for _, wrong := range []string{"linux", "debian", "darwin"} {
|
||||
got := problemsOf(t, bundleFor("go", wrong))
|
||||
if !strings.Contains(got, "and a system is") {
|
||||
t.Fatalf("%q was accepted as a system:\n%s", wrong, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// **And a language that runs anywhere must not name one.** A system that decides nothing reads as
|
||||
// though it did, which is the same fault as a restriction that restricts nothing (novox/hq ADR 0045).
|
||||
func TestAPortableBundleMayNotNameASystem(t *testing.T) {
|
||||
got := problemsOf(t, bundleFor("typescript", "arch"))
|
||||
if !strings.Contains(got, "runs anywhere") {
|
||||
t.Fatalf("a portable bundle was allowed to name a system:\n%s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAPortableBundleNamingNoSystemIsAccepted(t *testing.T) {
|
||||
if got := problemsOf(t, bundleFor("typescript", "")); got != "" {
|
||||
t.Fatalf("an ordinary bundle was refused:\n%s", got)
|
||||
}
|
||||
}
|
||||
|
||||
// One component, one artifact per system: the shape the mesh's own binaries are declared in, and the
|
||||
// reason the target is the artifact's rather than the recipe's.
|
||||
func TestOneArtifactPerSystemIsAccepted(t *testing.T) {
|
||||
m := Manifest{Module: "the-host", Build: &Build{Artifacts: []Artifact{
|
||||
{Name: "arch", Kind: ArtifactBundle, Language: "go", System: "arch"},
|
||||
{Name: "alpine", Kind: ArtifactBundle, Language: "go", System: "alpine"},
|
||||
{Name: "android", Kind: ArtifactBundle, Language: "go", System: "android"},
|
||||
}}}
|
||||
if got := problemsOf(t, m); got != "" {
|
||||
t.Fatalf("one artifact per system was refused:\n%s", got)
|
||||
}
|
||||
}
|
||||
@@ -124,6 +124,16 @@ type Rendering struct {
|
||||
// nothing on this node keeps them, or the mesh has no operator key.
|
||||
Kept *KeptExport
|
||||
|
||||
// OutwardLinks is the links this machine reported as facing outside it, which the filter is
|
||||
// written around (novox/hq ADR 0140). Empty means the machine has not said, and the mesh
|
||||
// composes no filter for it rather than writing a rule around a link with no name.
|
||||
OutwardLinks []string
|
||||
|
||||
// TunnelInterface is the interface the mesh's private network runs on, named here rather than
|
||||
// imported because the overlay package rests on this one. Traffic arriving on it is the mesh's,
|
||||
// not this machine's own guest, so the filter admits it only by a rule.
|
||||
TunnelInterface string
|
||||
|
||||
// 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
|
||||
// is the one that matters: a machine dials it to enrol, and a firewall derived only from
|
||||
@@ -345,7 +355,21 @@ func (r Resolution) compose(with Rendering, owner map[string]string) ([]map[stri
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
filtering := AsNftables(rules, with.Mesh, r.PublicDomain != "", with.Foundation)
|
||||
// **A machine that has not said which links face outside is sent no filter** (novox/hq ADR
|
||||
// 0140). The whole chain is written around those links: with none, the rule that lets this
|
||||
// machine's own guests keep working would name an empty set, which nftables refuses, and a rule
|
||||
// set that does not load is a machine filtering nothing while its unit reports success. Refused
|
||||
// here, where a person reads it, rather than on the machine — and the machine keeps the filter
|
||||
// it already has.
|
||||
if filters := r.filtersHere(); filters != "" && len(with.OutwardLinks) == 0 {
|
||||
return nil, fmt.Errorf(
|
||||
"%s cannot be sent a filter: it has not reported which of its links face outside, and "+
|
||||
"every rule in the chain is written around them. It reports that on each apply; "+
|
||||
"`node show %s` says whether it has. Until then %s is not sent, and the machine "+
|
||||
"keeps the filter it has", r.Node, r.Node, filters)
|
||||
}
|
||||
filtering := AsNftables(rules, with.Mesh, r.PublicDomain != "", with.Foundation,
|
||||
with.OutwardLinks, with.TunnelInterface)
|
||||
|
||||
var out []map[string]any
|
||||
for _, m := range r.Modules {
|
||||
@@ -852,6 +876,17 @@ func mapping(written string) (outer, inner int, address string, ok bool) {
|
||||
return outer, inner, strings.Join(parts[:len(parts)-2], ":"), true
|
||||
}
|
||||
|
||||
// filtersHere is the module on this node that loads the machine's packet filter, or empty when none
|
||||
// does. Named rather than counted: a refusal that says which module is one step from acted on.
|
||||
func (r Resolution) filtersHere() string {
|
||||
for _, m := range r.Modules {
|
||||
if m.Filtering != nil {
|
||||
return m.Module
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// Rules is the rule set this node's filter is derived from: every module's listens, what was
|
||||
// computed for this machine, and each module's per-node exposure. The same answer whether the node
|
||||
// is adopted or converged — the one loads it as a filter, the other declares it as openings.
|
||||
|
||||
@@ -230,7 +230,23 @@ const SSHPort = 22
|
||||
// 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
|
||||
// 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,
|
||||
outwardLinks []string, tunnel string) string {
|
||||
// The links that are not this machine's own: the ones facing outside, and the mesh's tunnel.
|
||||
// Traffic arriving on any of them is admitted only by a rule below; traffic arriving anywhere
|
||||
// else is this machine's own guest and is not something the mesh has a position on.
|
||||
//
|
||||
// The tunnel is named here deliberately. Treating it as "not outside" would make a port nothing
|
||||
// declares reachable from every machine in the mesh, which is the derivation abandoned.
|
||||
quoted := make([]string, 0, len(outwardLinks)+1)
|
||||
for _, link := range outwardLinks {
|
||||
quoted = append(quoted, fmt.Sprintf("%q", link))
|
||||
}
|
||||
if tunnel != "" {
|
||||
quoted = append(quoted, fmt.Sprintf("%q", tunnel))
|
||||
}
|
||||
inward := strings.Join(quoted, ", ")
|
||||
|
||||
var b strings.Builder
|
||||
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")
|
||||
@@ -252,6 +268,22 @@ func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int) str
|
||||
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")
|
||||
|
||||
// **What this machine's own guests must be able to ask it** (novox/hq ADR 0140). A guest gets
|
||||
// its address and its names from this machine, over the link 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.
|
||||
//
|
||||
// Asked for by the link it arrives on rather than by the address it comes from, for the reason
|
||||
// the forward chain below no longer names an address: a range describes one machine and goes
|
||||
// stale in silence. Anything arriving from outside, or over the tunnel, is not a guest of this
|
||||
// machine and asks through a port somebody declared, like everything else.
|
||||
if len(inward) > 0 {
|
||||
b.WriteString("\t\t# this machine's own guests asking it for an address and for names\n")
|
||||
b.WriteString(fmt.Sprintf("\t\tiifname != { %s } udp dport { 53, 67 } accept\n", inward))
|
||||
b.WriteString(fmt.Sprintf("\t\tiifname != { %s } tcp dport 53 accept\n", inward))
|
||||
}
|
||||
|
||||
// **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
|
||||
@@ -361,19 +393,36 @@ func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int) str
|
||||
// about the ports most worth protecting. Rehearsed on three machines: loading these rules
|
||||
// refused a port on the host and left a published container port reachable (novox/hq issue 047).
|
||||
//
|
||||
// The way through is the one the system being replaced already used: deny by default here, and
|
||||
// then explicitly allow the runtime's own networks, so containers keep working while everything
|
||||
// else has to be asked for.
|
||||
// **What it constrains is traffic arriving from OUTSIDE this machine, and nothing else**
|
||||
// (novox/hq ADR 0140).
|
||||
//
|
||||
// It used to deny everything here and then allow the machine's own containers back by naming
|
||||
// the address ranges they sit on — two ranges fixed in this file and the rest recorded per
|
||||
// machine. Every way of keeping that list correct failed. A constant describes one machine. A
|
||||
// recorded range goes stale in silence and cannot tell a network the mesh made from one a
|
||||
// predecessor left behind. Generating it from the modules would have put half this rule set on
|
||||
// the machine.
|
||||
//
|
||||
// The list should not exist, because the mesh has no position on a container reaching outward:
|
||||
// that is not a port opened to anybody. So traffic that did not arrive from outside is accepted
|
||||
// in one line, and what did arrive from outside is allowed only where a rule below admits it.
|
||||
//
|
||||
// The tunnel is not "not outside". Accepting everything off it would make a port nothing
|
||||
// declares reachable from any machine in the mesh, which is the derivation abandoned — so it is
|
||||
// named here beside the outward links, and traffic arriving on it meets the rules below like
|
||||
// anything else.
|
||||
b.WriteString("\tchain forward {\n")
|
||||
b.WriteString("\t\ttype filter hook forward priority filter; policy drop;\n")
|
||||
b.WriteString("\t\tct state established,related accept\n")
|
||||
b.WriteString("\t\tct state invalid drop\n")
|
||||
b.WriteString("\n")
|
||||
// What the container runtime created. Without these, denying by default stops every container
|
||||
// on the machine — which is exactly the failure the absent chain was avoiding, avoided properly.
|
||||
for _, network := range runtimeNetworks {
|
||||
b.WriteString(fmt.Sprintf("\t\t# %s\n", network.why))
|
||||
b.WriteString(fmt.Sprintf("\t\tip saddr %s accept\n", network.cidr))
|
||||
// Only when there is a link to name. An empty set is a line nftables refuses, and a rule set
|
||||
// that does not load is a machine filtering nothing while its unit reports success — so the
|
||||
// chain denies rather than renders nonsense. Composing a declaration for a machine that has
|
||||
// named none is refused upstream, so this is a floor and not a path anything travels.
|
||||
if inward != "" {
|
||||
b.WriteString("\t\t# this machine's own guests reaching outward: not a port opened to anybody\n")
|
||||
b.WriteString(fmt.Sprintf("\t\tiifname != { %s } accept\n", inward))
|
||||
}
|
||||
|
||||
if len(rules) > 0 {
|
||||
@@ -430,18 +479,6 @@ func AsNftables(rules []Rule, mesh []string, outward bool, foundation []int) str
|
||||
return b.String()
|
||||
}
|
||||
|
||||
// runtimeNetworks are the container runtime's own networks, which must keep working when the
|
||||
// forward chain denies by default.
|
||||
//
|
||||
// Taken from what the system being replaced allows, which has been carrying this machine's traffic
|
||||
// for months: the runtime's bridge range and the range its compose files are given. A machine whose
|
||||
// runtime is configured with something else needs this to say so — which is a thing the mesh cannot
|
||||
// derive and a reason this list is named here rather than computed.
|
||||
var runtimeNetworks = []struct{ cidr, why string }{
|
||||
{"172.16.0.0/12", "the container runtime's bridge networks"},
|
||||
{"192.168.128.0/17", "the networks its compose files are given"},
|
||||
}
|
||||
|
||||
// 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
|
||||
|
||||
@@ -19,7 +19,7 @@ func TestTheBrokersPortIsOpenedThoughNoModuleDeclaresIt(t *testing.T) {
|
||||
// A machine on the private network, with one ordinary module rule, and nothing that mentions
|
||||
// the broker — which is every machine.
|
||||
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, "mesh0")
|
||||
|
||||
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)
|
||||
@@ -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
|
||||
// a panic. A control plane in that state cannot issue tokens either, which is where it surfaces.
|
||||
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, "mesh0")
|
||||
if !strings.Contains(out, "table inet mesh") {
|
||||
t.Fatalf("no ruleset at all:\n%s", out)
|
||||
}
|
||||
|
||||
@@ -79,7 +79,7 @@ func TestTwoModulesWantingOnePortAreBothNamed(t *testing.T) {
|
||||
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.
|
||||
nft := AsNftables(rules, nil, false, nil)
|
||||
nft := AsNftables(rules, nil, false, nil, nil, "mesh0")
|
||||
if !strings.Contains(nft, "web") || !strings.Contains(nft, "board") {
|
||||
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) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{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, "mesh0")
|
||||
// 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.
|
||||
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,
|
||||
// including the ones the container runtime writes for its bridges.
|
||||
func TestReloadingReplacesOnlyTheMeshsOwnRules(t *testing.T) {
|
||||
nft := AsNftables(nil, nil, false, nil)
|
||||
nft := AsNftables(nil, nil, false, nil, nil, "mesh0")
|
||||
if strings.Contains(nft, "flush ruleset") {
|
||||
t.Fatalf("loading the rule set empties every table on the machine:\n%s", nft)
|
||||
}
|
||||
@@ -160,22 +160,122 @@ func TestReloadingReplacesOnlyTheMeshsOwnRules(t *testing.T) {
|
||||
// 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.
|
||||
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, "mesh0")
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
// And containers keep working, which is the whole reason the chain was left out before.
|
||||
func TestTheRuntimesOwnNetworksKeepWorking(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil)
|
||||
for _, network := range []string{"172.16.0.0/12", "192.168.128.0/17"} {
|
||||
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)
|
||||
// And this machine's own guests keep working, which is the whole reason the chain was left out
|
||||
// before — by not being mentioned (novox/hq ADR 0140).
|
||||
//
|
||||
// It used to be done by naming the address ranges they sit on: two fixed here and the rest recorded
|
||||
// per machine. That list broke a workstation's containers at a flip and could not be made correct,
|
||||
// because a range describes one machine and cannot tell a network the mesh made from one a
|
||||
// predecessor left behind. What replaced it is a single line about the links traffic arrives on.
|
||||
func TestThisMachinesOwnGuestsKeepWorkingWithoutBeingNamed(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, []string{"eth0"}, "mesh0")
|
||||
if !strings.Contains(nft, `iifname != { "eth0", "mesh0" } accept`) {
|
||||
t.Fatalf("what did not arrive from outside is not accepted, so this machine's own guests "+
|
||||
"reach nothing:\n%s", nft)
|
||||
}
|
||||
}
|
||||
|
||||
// No address of a machine's own networks appears anywhere in a rendered filter.
|
||||
//
|
||||
// This is the assertion that fails against the previous behaviour, and it is why it is written on
|
||||
// the text rather than on an outcome: the two ranges were a constant in this file, so nothing but
|
||||
// reading the output catches one creeping back in.
|
||||
func TestNoNetworkOfTheMachinesOwnIsNamed(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, []string{"eth0"}, "mesh0")
|
||||
for _, gone := range []string{"172.16.0.0/12", "192.168.128.0/17", "saddr 192.168", "saddr 172."} {
|
||||
if strings.Contains(nft, gone) {
|
||||
t.Fatalf("%q is named, and a range describes one machine and goes stale in silence:\n%s",
|
||||
gone, nft)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// **The tunnel is constrained, not treated as inside.**
|
||||
//
|
||||
// Accepting everything arriving over the private network would make a port nothing declares
|
||||
// reachable from every machine in the mesh — the derivation abandoned, and a rule that reads as a
|
||||
// restriction while restricting nothing. So the tunnel is named beside the outward links, and
|
||||
// traffic arriving on it meets the declared rules like anything else.
|
||||
func TestTheTunnelIsConstrainedLikeAnOutwardLink(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, []string{"eth0"}, "mesh0")
|
||||
line := `iifname != { "eth0", "mesh0" } accept`
|
||||
if !strings.Contains(nft, line) {
|
||||
t.Fatalf("the tunnel is not constrained, so an undeclared port is reachable from any "+
|
||||
"machine in the mesh:\n%s", nft)
|
||||
}
|
||||
}
|
||||
|
||||
// A machine with two links facing outside has both constrained. Asserted on the one line, because a
|
||||
// rule covering one and not the other would leave a machine filtering half of what reaches it.
|
||||
func TestEveryOutwardLinkIsConstrained(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, []string{"eth0", "wlan0"}, "mesh0")
|
||||
if !strings.Contains(nft, `iifname != { "eth0", "wlan0", "mesh0" } accept`) {
|
||||
t.Fatalf("not every outward link is constrained:\n%s", nft)
|
||||
}
|
||||
}
|
||||
|
||||
// A guest asks its host for an address and for names, and those two arrive at the input chain. Asked
|
||||
// for by the link they arrive on, so a resolver bound anywhere but an outward link keeps answering.
|
||||
func TestGuestsMayAskTheirHostForAnAddressAndNames(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, []string{"eth0"}, "mesh0")
|
||||
for _, want := range []string{
|
||||
`iifname != { "eth0", "mesh0" } udp dport { 53, 67 } accept`,
|
||||
`iifname != { "eth0", "mesh0" } tcp dport 53 accept`,
|
||||
} {
|
||||
if !strings.Contains(nft, want) {
|
||||
t.Fatalf("a guest cannot ask its host for an address or a name, which is not a closed "+
|
||||
"port but a network that does not work:\n%s", nft)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// With no link named at all the chain denies rather than rendering an empty set, which nftables
|
||||
// refuses — and a rule set that does not load is a machine filtering nothing while its unit reports
|
||||
// success. Composing a declaration for such a machine is refused upstream; this is the floor.
|
||||
func TestNoLinkNamedRendersNoCatchAllRatherThanAnEmptySet(t *testing.T) {
|
||||
nft := AsNftables(nil, []string{"198.51.100.2"}, false, nil, nil, "")
|
||||
if strings.Contains(nft, "{ }") || strings.Contains(nft, "iifname != {}") {
|
||||
t.Fatalf("an empty set is rendered, which nftables refuses:\n%s", nft)
|
||||
}
|
||||
if !strings.Contains(nft, "hook forward priority filter; policy drop") {
|
||||
t.Fatalf("the forward chain does not deny:\n%s", nft)
|
||||
}
|
||||
}
|
||||
|
||||
// A machine that has not said which links face outside is sent no filter, and the refusal names the
|
||||
// module that would have loaded it so the reader knows what is being withheld.
|
||||
func TestAMachineThatNamedNoOutwardLinkIsSentNoFilter(t *testing.T) {
|
||||
r := Resolution{Node: "anchor", Modules: []Manifest{
|
||||
{Module: "nftables", Filtering: &Filtering{Into: "/etc/mesh/filter.nft"}},
|
||||
{Module: "web", Listens: []Listening{{Port: 443, From: FromEverywhere}}},
|
||||
}}
|
||||
_, err := r.Declaration(Rendering{Mesh: []string{"198.51.100.2"}, TunnelInterface: "mesh0"})
|
||||
if err == nil {
|
||||
t.Fatal("a machine that named no outward link was sent a filter written around none")
|
||||
}
|
||||
for _, want := range []string{"anchor", "nftables", "face outside"} {
|
||||
if !strings.Contains(err.Error(), want) {
|
||||
t.Fatalf("the refusal does not say %q: %v", want, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// And a machine that names none but loads no filter is not refused: there is nothing to write.
|
||||
func TestAMachineWithNoFilterModuleIsNotRefused(t *testing.T) {
|
||||
r := Resolution{Node: "anchor", Modules: []Manifest{
|
||||
{Module: "web", Listens: []Listening{{Port: 443, From: FromEverywhere}}},
|
||||
}}
|
||||
if _, err := r.Declaration(Rendering{Mesh: []string{"198.51.100.2"}}); err != nil {
|
||||
t.Fatalf("a machine that loads no filter was refused one: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A published port is matched by what the client asked for, not by where the packet ends up.
|
||||
//
|
||||
// The runtime rewrites the destination before this chain sees it, so a rule naming the published
|
||||
@@ -183,7 +283,7 @@ func TestTheRuntimesOwnNetworksKeepWorking(t *testing.T) {
|
||||
func TestAPublishedPortIsMatchedByWhatWasAskedFor(t *testing.T) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{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, "mesh0")
|
||||
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)
|
||||
}
|
||||
@@ -193,7 +293,7 @@ func TestAPublishedPortIsMatchedByWhatWasAskedFor(t *testing.T) {
|
||||
func TestAMeshScopedPortIsMeshScopedWhenForwarded(t *testing.T) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{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, "mesh0")
|
||||
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)
|
||||
}
|
||||
@@ -203,7 +303,7 @@ func TestAMeshScopedPortIsMeshScopedWhenForwarded(t *testing.T) {
|
||||
func TestFromTheMeshIsTheNodesTheMeshKnows(t *testing.T) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{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, "mesh0")
|
||||
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)
|
||||
}
|
||||
@@ -213,7 +313,7 @@ func TestFromTheMeshIsTheNodesTheMeshKnows(t *testing.T) {
|
||||
func TestAMeshPortOnANodeWithNoMeshIsClosedAndSaysSo(t *testing.T) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{Module: "store", Listens: []Listening{{Port: 5432, From: FromMesh}}},
|
||||
}}, nil), nil, false, nil)
|
||||
}}, nil), nil, false, nil, nil, "mesh0")
|
||||
if strings.Contains(nft, "dport 5432 accept") {
|
||||
t.Fatalf("a port meant for the mesh was opened to everything:\n%s", nft)
|
||||
}
|
||||
@@ -226,7 +326,7 @@ func TestAMeshPortOnANodeWithNoMeshIsClosedAndSaysSo(t *testing.T) {
|
||||
func TestAMachineScopedPortIsNotOpened(t *testing.T) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{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, "mesh0")
|
||||
if strings.Contains(nft, "dport 6379 accept") {
|
||||
t.Fatalf("a port for this machine only was opened to the network:\n%s", nft)
|
||||
}
|
||||
@@ -238,7 +338,8 @@ func TestTheModuleAskingForTheRuleSetGetsEveryModulesPorts(t *testing.T) {
|
||||
{Module: "firewall", Filtering: &Filtering{Into: "/etc/mesh/filter.nft"}},
|
||||
{Module: "web", Listens: []Listening{{Port: 443, From: FromEverywhere}}},
|
||||
}}
|
||||
out, err := r.Declaration(Rendering{Mesh: []string{"198.51.100.2"}})
|
||||
out, err := r.Declaration(Rendering{Mesh: []string{"198.51.100.2"},
|
||||
OutwardLinks: []string{"eth0"}, TunnelInterface: "mesh0"})
|
||||
if err != nil {
|
||||
t.Fatalf("declaration: %v", err)
|
||||
}
|
||||
@@ -268,7 +369,7 @@ func TestAskingForTheRuleSetWithNowhereToPutItIsRefused(t *testing.T) {
|
||||
func TestAMeshOnBothAddressFamiliesRendersBoth(t *testing.T) {
|
||||
nft := AsNftables(mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{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, "mesh0")
|
||||
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)
|
||||
}
|
||||
@@ -296,7 +397,8 @@ func TestWhatTheMeshComputesIsAppliedBeforeWhatTheModuleDeclared(t *testing.T) {
|
||||
"restart-on": []any{"filtering"}},
|
||||
},
|
||||
}}}
|
||||
out, err := r.Declaration(Rendering{Mesh: []string{"198.51.100.2"}})
|
||||
out, err := r.Declaration(Rendering{Mesh: []string{"198.51.100.2"},
|
||||
OutwardLinks: []string{"eth0"}, TunnelInterface: "mesh0"})
|
||||
if err != nil {
|
||||
t.Fatalf("declaration: %v", err)
|
||||
}
|
||||
@@ -672,7 +774,7 @@ func TestExposureRefusesAPortNotListenedOnAndABadSource(t *testing.T) {
|
||||
// loading the rules lives on conntrack until it drops, and then the machine is reached from a
|
||||
// rescue console (novox/hq issue 047).
|
||||
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, "mesh0")
|
||||
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)
|
||||
}
|
||||
@@ -685,7 +787,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
|
||||
// private network is the thing that broke.
|
||||
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, "mesh0")
|
||||
if !strings.Contains(nft, "\t\ttcp dport 22 accept") {
|
||||
t.Fatalf("a machine reachable from outside does not answer ssh there:\n%s", nft)
|
||||
}
|
||||
@@ -697,7 +799,7 @@ func TestSSHIsOpenFromOutsideOnAMachineThatFacesIt(t *testing.T) {
|
||||
// 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.
|
||||
func TestSSHIsNeverLeftWithoutARule(t *testing.T) {
|
||||
nft := AsNftables(nil, nil, false, nil)
|
||||
nft := AsNftables(nil, nil, false, nil, nil, "mesh0")
|
||||
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)
|
||||
}
|
||||
|
||||
@@ -571,6 +571,21 @@ type Artifact struct {
|
||||
// image built from this same module's own repository, the same as every other artifact.
|
||||
Context *ArtifactContext `json:"context,omitempty"`
|
||||
|
||||
// System is the operating system this artifact is compiled for, for a bundle whose output is a
|
||||
// binary rather than portable code (novox/hq ADR 0142).
|
||||
//
|
||||
// **Named by the artifact, not by the recipe.** A toolchain deliberately accepts nothing from
|
||||
// the module — anything a module could override there it would be writing a Dockerfile to
|
||||
// override — and yet a compiled binary is per operating system, pinned at link time so a host
|
||||
// refuses to touch a machine it was not built for (novox/hq ADR 0005). The way out is that the
|
||||
// target is a property of the artifact: one artifact declared per system, one build each, and
|
||||
// the recipe stays the mesh's.
|
||||
//
|
||||
// Empty for a bundle whose output runs anywhere, which is every interpreted language, and for
|
||||
// every other kind. A bundle in a language that compiles to a binary must say one, because
|
||||
// "compiled for whatever the build machine happened to be" is the fault this exists to prevent.
|
||||
System string `json:"system,omitempty"`
|
||||
|
||||
// Language is what this module's code is written in, for a bundle.
|
||||
//
|
||||
// **Declared, never guessed.** Inferring it from what files happen to be present makes a
|
||||
|
||||
@@ -12,5 +12,5 @@ func TestPrintRehearsalRuleset(t *testing.T) {
|
||||
rules := mustFilter(t, Resolution{Modules: []Manifest{
|
||||
{Module: "pub", Listens: []Listening{{Port: 8099, From: FromMesh, Why: "the thing it serves"}}},
|
||||
}}, 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, "mesh0"))
|
||||
}
|
||||
|
||||
@@ -48,7 +48,11 @@ type Seat struct {
|
||||
//
|
||||
// In the order a person reads it: the mesh's own, then a node's.
|
||||
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"},
|
||||
// **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
|
||||
|
||||
@@ -30,12 +30,12 @@ func TestRefusedAndFailedAreDifferentSituations(t *testing.T) {
|
||||
refuser := nodeNamed(t, inv, "refuser")
|
||||
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",
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.RecordDoing(ctx, failer, Doing{
|
||||
if _, err := inv.RecordDoing(ctx, failer, Doing{
|
||||
Outcome: OutcomeFailed,
|
||||
Failed: []FailedResource{{ID: "svc", Error: "unit not found"}},
|
||||
Applied: 4,
|
||||
@@ -70,7 +70,7 @@ func TestAMachineDoingWhatItWasToldIsNotOnTheList(t *testing.T) {
|
||||
inv := fresh(t)
|
||||
ctx := context.Background()
|
||||
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)
|
||||
}
|
||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||
@@ -97,12 +97,12 @@ func TestTheLastReportReplacesTheOneBefore(t *testing.T) {
|
||||
inv := fresh(t)
|
||||
ctx := context.Background()
|
||||
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"}},
|
||||
}); err != nil {
|
||||
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)
|
||||
}
|
||||
wrong, err := inv.NotDoingWhatTheyWereTold(ctx)
|
||||
@@ -141,7 +141,7 @@ func TestWhatANodeSaidGoesWhenTheNodeDoes(t *testing.T) {
|
||||
inv := fresh(t)
|
||||
ctx := context.Background()
|
||||
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)
|
||||
}
|
||||
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")
|
||||
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)
|
||||
}
|
||||
first, _, err := inv.DoingOf(ctx, "looping")
|
||||
@@ -269,7 +269,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -287,7 +287,7 @@ func TestTheSameFailureReportedAgainIsCountedNotRestarted(t *testing.T) {
|
||||
// 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.
|
||||
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)
|
||||
}
|
||||
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.
|
||||
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)
|
||||
}
|
||||
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.
|
||||
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)
|
||||
}
|
||||
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.
|
||||
if err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
if _, err := inv.RecordDoing(ctx, id, same); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
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;
|
||||
@@ -0,0 +1,26 @@
|
||||
-- Which of a machine's links face outside it, replacing the networks it was told to say it routes.
|
||||
--
|
||||
-- novox/hq ADR 0140, superseding 0137 and 0139. The derived filter blocked everything passing
|
||||
-- through a machine and then allowed the machine's own containers back by naming the address ranges
|
||||
-- they sit on: two ranges fixed in the controller's source, the rest recorded by 0043's column.
|
||||
--
|
||||
-- Every route to a correct list fails. A constant describes one machine. A recorded range goes stale
|
||||
-- in silence, and cannot tell a network the mesh made from one a predecessor left behind — measured
|
||||
-- on the control-node, where six ranges fall outside the constants and two of the six belong to
|
||||
-- services the mesh does not run. Generating the list from the modules put half the rule set on the
|
||||
-- machine.
|
||||
--
|
||||
-- The list should not exist, because the mesh has no position on a container reaching outward: that
|
||||
-- is not a port opened to anybody. The filter constrains what arrives from OUTSIDE the machine and
|
||||
-- says nothing about what did not, which needs one fact instead of a list — which links "outside"
|
||||
-- arrives on.
|
||||
--
|
||||
-- Reported by the machine on every apply, never recorded by hand, so it cannot go stale. Null for a
|
||||
-- machine that has not reported yet; the mesh composes no filter for such a machine and leaves the
|
||||
-- one it has, because a rule written around a link with no name is a rule set that does not load.
|
||||
alter table node add column outward_links jsonb;
|
||||
|
||||
-- What 0043 recorded is not migrated into it. The ranges answered a question that no longer exists,
|
||||
-- and every machine that named one keeps working without it: the traffic those ranges allowed is now
|
||||
-- allowed by not having arrived from outside.
|
||||
alter table node drop column routed_networks;
|
||||
+84
-18
@@ -518,6 +518,61 @@ func (i *Inventory) PublicDomainOf(ctx context.Context, name string) (string, er
|
||||
return *domain, nil
|
||||
}
|
||||
|
||||
// RecordOutwardLinks keeps the links a machine reported as facing outside it.
|
||||
//
|
||||
// A reported fact, not a setting (novox/hq ADR 0140). It replaces the networks a machine used to be
|
||||
// told to say it routes: the filter blocked everything passing through and then allowed the machine's
|
||||
// own containers back by naming their address ranges, and every way of keeping that list correct
|
||||
// failed — a constant describes one machine, and a recorded range goes stale in silence. The filter
|
||||
// now constrains what arrives from outside and says nothing about what did not, and the one thing it
|
||||
// needs is which links "outside" arrives on. The machine reads that from its own routing table on
|
||||
// every apply, so it cannot go stale and nobody types it.
|
||||
//
|
||||
// An empty list clears it, which is what a machine with no route off itself reports. The mesh then
|
||||
// composes no filter for that machine at all.
|
||||
func (i *Inventory) RecordOutwardLinks(ctx context.Context, id string, links []string) error {
|
||||
var kept []string
|
||||
for _, name := range links {
|
||||
if name = strings.TrimSpace(name); name != "" {
|
||||
kept = append(kept, name)
|
||||
}
|
||||
}
|
||||
if len(kept) == 0 {
|
||||
_, err := i.store.Pool().Exec(ctx,
|
||||
`update node set outward_links = null where id = $1`, id)
|
||||
return err
|
||||
}
|
||||
body, err := json.Marshal(kept)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = i.store.Pool().Exec(ctx,
|
||||
`update node set outward_links = $2 where id = $1`, id, string(body))
|
||||
return err
|
||||
}
|
||||
|
||||
// OutwardLinksOf is the links a machine reported as facing outside it, empty when it has reported
|
||||
// none — which is a machine the mesh composes no filter for.
|
||||
func (i *Inventory) OutwardLinksOf(ctx context.Context, name string) ([]string, error) {
|
||||
var body []byte
|
||||
err := i.store.Pool().QueryRow(ctx,
|
||||
`select outward_links 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 links []string
|
||||
if err := json.Unmarshal(body, &links); err != nil {
|
||||
return nil, fmt.Errorf("the outward links recorded for %s are not a list: %w", name, err)
|
||||
}
|
||||
return links, nil
|
||||
}
|
||||
|
||||
// RecordOverlayKey keeps the public half a node generated.
|
||||
func (i *Inventory) RecordOverlayKey(ctx context.Context, node, key string) error {
|
||||
if strings.TrimSpace(key) == "" {
|
||||
@@ -670,31 +725,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
|
||||
// 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.
|
||||
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)
|
||||
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
|
||||
times := 0
|
||||
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()
|
||||
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
|
||||
}
|
||||
}
|
||||
@@ -707,7 +770,10 @@ func (i *Inventory) RecordDoing(ctx context.Context, node string, d Doing) error
|
||||
declared = excluded.declared,
|
||||
failing_since = excluded.failing_since, failures = excluded.failures`,
|
||||
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.
|
||||
|
||||
@@ -30,6 +30,11 @@ type Bus interface {
|
||||
// (design 29 §4, the *state* shape).
|
||||
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
|
||||
// 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.
|
||||
@@ -81,6 +86,13 @@ func EventSubject(source, key string) string {
|
||||
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
|
||||
// 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.
|
||||
@@ -112,6 +124,27 @@ func (b OverNATS) PublishEvent(ctx context.Context, key, source, node string, bo
|
||||
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 {
|
||||
_, err := b.JS.Publish(DeclareSubject(node), body, nats.Context(ctx))
|
||||
if err != nil {
|
||||
|
||||
+27
-13
@@ -265,7 +265,7 @@ func (e Enrolment) Outstanding(ctx context.Context, node string) (string, error)
|
||||
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
|
||||
// another attempt rather than acknowledged and lost (novox/hq issue 082).
|
||||
defer func() {
|
||||
@@ -274,11 +274,11 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
}
|
||||
}()
|
||||
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)
|
||||
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
|
||||
@@ -298,7 +298,18 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
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 {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
// Which of its links face outside (novox/hq ADR 0140), whenever it says so. Recorded on every
|
||||
// report that carries it, adopted or converged, because the filter the mesh composes is written
|
||||
// around it — and never cleared by a report that carries none, which is every bare word that the
|
||||
// node is there. A machine whose routing table it could not read reports nothing rather than
|
||||
// guessing, and keeps whatever it last said; a machine with genuinely no route off itself is one
|
||||
// the mesh composes no filter for at all.
|
||||
if len(report.Outward) > 0 {
|
||||
if err := e.Inventory.RecordOutwardLinks(ctx, node.ID, report.Outward); err != nil {
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
// What it says about the tunnel it carried (novox/hq ADR 0105), whenever it says it.
|
||||
@@ -308,7 +319,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
Peers: report.Tunnel.Peers, State: report.Tunnel.State, Note: report.Tunnel.Note,
|
||||
Kept: report.Tunnel.Kept,
|
||||
}); err != nil {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
// A node taking a found tunnel's key after enrolment (novox/hq ADR 0105). Verified against the
|
||||
@@ -317,9 +328,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.
|
||||
if report.Rekey != 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
|
||||
@@ -334,7 +345,7 @@ func (e Enrolment) Heard(ctx context.Context, report Report) (err error) {
|
||||
if 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
|
||||
// last_seen and the reason went to a log line, so "which machine is not doing what it was
|
||||
@@ -361,17 +372,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
|
||||
// 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 {
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
if err := e.Inventory.RecordDoing(ctx, node.ID, doing); err != nil {
|
||||
return err
|
||||
// **Whether this is news** is the store's answer: it holds the previous report, and a machine
|
||||
// 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
|
||||
// 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.
|
||||
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
|
||||
}
|
||||
|
||||
// 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
|
||||
// it in the module graph; nothing else need care.
|
||||
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 {
|
||||
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)
|
||||
}
|
||||
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 {
|
||||
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},
|
||||
}); err != nil {
|
||||
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.
|
||||
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)
|
||||
}
|
||||
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 {
|
||||
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"},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -200,13 +200,13 @@ func TestWhatAnAdoptedNodeHoldsIsKeptAndAnAliveWordDoesNotWipeIt(t *testing.T) {
|
||||
}
|
||||
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)
|
||||
}
|
||||
check("after an alive word")
|
||||
|
||||
// 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 {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@@ -170,6 +170,20 @@ type Report struct {
|
||||
// Firewall is the firewall found on the machine — "ufw" or "none" — and empty on a node that
|
||||
// was never asked, which is every converged one.
|
||||
Firewall string `json:"firewall,omitempty"`
|
||||
|
||||
// Outward is the links on this machine that face outside it — the ones carrying a default route
|
||||
// (novox/hq ADR 0140). Every node reports it, adopted or converged, because the filter the mesh
|
||||
// composes for it is written around these and nothing else.
|
||||
//
|
||||
// **It replaces a list of addresses.** The filter used to block everything passing through the
|
||||
// machine and then allow the machine's own containers back by naming the ranges they sit on. A
|
||||
// range describes one machine and goes stale in silence; the link carrying the default route is
|
||||
// read afresh on every report and does not change when a module is added or removed.
|
||||
//
|
||||
// Empty means the machine has not said. The mesh composes no filter for such a machine and
|
||||
// leaves the one it has: a rule written around a link with no name is a rule set that does not
|
||||
// load, and that is a machine filtering nothing while its unit reports success.
|
||||
Outward []string `json:"outward,omitempty"`
|
||||
// Reachable is what can be reached on the machine now: every listening socket and every
|
||||
// published container port. Only an adopted node reports it; it is what converging previews.
|
||||
Reachable []Reach `json:"reachable,omitempty"`
|
||||
|
||||
@@ -86,14 +86,15 @@ type counted struct {
|
||||
heard []Report
|
||||
}
|
||||
|
||||
func (c *counted) Heard(_ context.Context, r Report) error {
|
||||
func (c *counted) Heard(_ context.Context, r Report) (bool, error) {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
if c.err != nil {
|
||||
return c.err
|
||||
return false, c.err
|
||||
}
|
||||
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) {
|
||||
@@ -207,11 +208,11 @@ type sentAndHeardSafely struct {
|
||||
heard []Report
|
||||
}
|
||||
|
||||
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) error {
|
||||
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) (bool, error) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.heard = append(s.heard, r)
|
||||
return nil
|
||||
return true, nil
|
||||
}
|
||||
|
||||
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.
|
||||
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.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)
|
||||
}
|
||||
placed, err := e.Inventory.Overlays(ctx)
|
||||
@@ -77,7 +77,7 @@ func TestASignedRekeyMovesTheHubOntoItsTunnel(t *testing.T) {
|
||||
_ = hub
|
||||
|
||||
// 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") {
|
||||
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.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") {
|
||||
t.Fatalf("a rekey signed by a stranger was accepted: %v", err)
|
||||
}
|
||||
@@ -111,7 +111,7 @@ func TestARekeySignedByAnotherKeyIsRefusedAndChangesNothing(t *testing.T) {
|
||||
other := theTunnel()
|
||||
other.Port = 51820
|
||||
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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,12 +9,12 @@ import (
|
||||
|
||||
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.
|
||||
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 {
|
||||
return Report{Node: node, Declared: declared, Applied: []string{"store"}}
|
||||
|
||||
+64
-5
@@ -29,7 +29,12 @@ type Enroller interface {
|
||||
// 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.
|
||||
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.
|
||||
@@ -257,7 +262,7 @@ func (s *Server) heartbeat(m Control) {
|
||||
return
|
||||
}
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -291,7 +296,7 @@ func (s *Server) reported(ctx context.Context, m Control) {
|
||||
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) {
|
||||
case Hold:
|
||||
// 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.
|
||||
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 {
|
||||
@@ -334,7 +346,8 @@ func (s *Server) reported(ctx context.Context, m Control) {
|
||||
// whenever it arrives, which is the behaviour the mesh has had all along.
|
||||
func staleAgainst(report Report) string {
|
||||
if report.Rekey != nil || report.Tunnel != nil || len(report.Held) > 0 ||
|
||||
report.Firewall != "" || len(report.Reachable) > 0 || len(report.Carried) > 0 {
|
||||
report.Firewall != "" || len(report.Reachable) > 0 || len(report.Carried) > 0 ||
|
||||
len(report.Outward) > 0 {
|
||||
return ""
|
||||
}
|
||||
return report.Declared
|
||||
@@ -460,7 +473,19 @@ func (s *Server) catchingUp(ctx context.Context, m Control) {
|
||||
sent := 0
|
||||
for _, a := range announcements {
|
||||
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
|
||||
// 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",
|
||||
@@ -551,3 +576,37 @@ func (s *Server) sourceMoved(ctx context.Context, m Control) {
|
||||
}
|
||||
_ = 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
|
||||
}
|
||||
|
||||
func (s *sentAndHeard) Heard(_ context.Context, r Report) error {
|
||||
func (s *sentAndHeard) Heard(_ context.Context, r Report) (bool, error) {
|
||||
if s.err != nil {
|
||||
return s.err
|
||||
return false, s.err
|
||||
}
|
||||
s.heard = append(s.heard, r)
|
||||
return nil
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (s *sentAndHeard) Outstanding(context.Context, string) (string, error) { return s.sent, nil }
|
||||
|
||||
Reference in New Issue
Block a user