Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b0d11c2439 | ||
|
|
155689672c | ||
|
|
e82789a322 | ||
|
|
99562947f3 | ||
|
|
c171c64d3a | ||
|
|
0dc5515099 | ||
|
|
58c0e715c7 | ||
|
|
c5b229f95d | ||
|
|
a9c49724b5 | ||
|
|
296ec5ece6 | ||
|
|
45b9a507a1 | ||
|
|
5c410aaf54 | ||
|
|
2478b5f127 | ||
|
|
c82e933268 | ||
|
|
e212da6bd2 | ||
|
|
673912ccdc | ||
|
|
68a193d6f0 | ||
|
|
4fba884d47 | ||
|
|
15f0dabf32 | ||
|
|
4ef41ad5a0 | ||
|
|
d749de989c |
+18
-7
@@ -33,6 +33,7 @@ import (
|
||||
"github.com/novox/mesh-host/internal/identity"
|
||||
"github.com/novox/mesh-host/internal/inventory"
|
||||
"github.com/novox/mesh-host/internal/link"
|
||||
"github.com/novox/mesh-host/internal/outward"
|
||||
"github.com/novox/mesh-host/internal/profile"
|
||||
"github.com/novox/mesh-host/internal/reachable"
|
||||
"github.com/novox/mesh-host/internal/store"
|
||||
@@ -962,9 +963,14 @@ func runLink(ctx context.Context, opts options) error {
|
||||
return nil // asked to stop while waiting
|
||||
}
|
||||
|
||||
fmt.Printf("node %s, linking to %s\n", mine.Node, mine.Membership.Broker)
|
||||
|
||||
// **A membership delivered while this host was not running is adopted before the first dial.**
|
||||
// The ordinary path is a declaration, read after it applies; the rescue path is an operator
|
||||
// writing the file by hand on a machine no bus can reach — rotated while it held the old
|
||||
// password, say — and restarting the host (design 28, task 5.2). Same file, same check.
|
||||
say := func(line string) { fmt.Println(line) }
|
||||
adoptDeliveredMembership(identity.Path(opts.state), &mine, say)
|
||||
|
||||
fmt.Printf("node %s, linking to %s\n", mine.Node, mine.Membership.Broker)
|
||||
|
||||
// One scheduler for the life of the process, re-established from each applied declaration
|
||||
// (novox/hq ADR 0053). It fires scheduled steps on their cadence, surviving across applies and
|
||||
@@ -1258,6 +1264,15 @@ func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.D
|
||||
}
|
||||
|
||||
report := link.Report{Carried: carriedPorts(updated), Declared: digestOf(raw)}
|
||||
// Which of this machine's links face outside, for the filter the mesh writes around them
|
||||
// (novox/hq ADR 0140). Reported whatever the node's mode: a converged node's filter needs it,
|
||||
// and an adopted one becomes converged without a further round trip. A machine that cannot read
|
||||
// its own routing table says nothing rather than guessing, and is sent no filter.
|
||||
if links, err := outward.Links(""); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "mesh-host: applied, and could not read which links face outside: %v\n", err)
|
||||
} else {
|
||||
report.Outward = links
|
||||
}
|
||||
// What this node found and holds, its firewall, and what is reachable on it — so an adopted
|
||||
// node never reads as converged (novox/hq ADR 0100).
|
||||
for _, h := range updated.Held {
|
||||
@@ -1428,10 +1443,6 @@ func adoptDeliveredMembership(identityPath string, mine *identity.Identity, say
|
||||
say(fmt.Sprintf("a membership for another bus was delivered and could not be saved: %v", err))
|
||||
return
|
||||
}
|
||||
transport := next.Transport
|
||||
if transport == "" {
|
||||
transport = "the current"
|
||||
}
|
||||
say(fmt.Sprintf("moving to %s bus at %s — restarting to dial it", transport, next.Broker))
|
||||
say(fmt.Sprintf("moving to the %s bus at %s — restarting to dial it", next.Transport, next.Broker))
|
||||
os.Exit(0)
|
||||
}
|
||||
|
||||
@@ -152,7 +152,7 @@
|
||||
"type": "file",
|
||||
"path": "/var/lib/mesh-bus-conf/accounts.conf",
|
||||
"mode": "0600",
|
||||
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.API.>\", \"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"_INBOX.enrol.>\", \"mesh.control.>\", \"mesh.node.>\", \"mesh.seat.mesh-build-machine.accept.>\"] }\n subscribe: { allow: [\"$JS.API.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.seat.mesh-build-machine.event.built\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
|
||||
"content": "// The first user list, carried by the installer because at genesis there is no mesh to\n// compose one. A bootstrap credential, rotated with the store's and replaced by the\n// controller's own composition from its first start onward.\naccounts {\n MESH {\n jetstream: enabled\n users = [\n { user: \"controller\", password: \"$2a$10$AHqJgOifIVbU41KmATiMhuXFs8xa7Wl2HuN4UVBCXdN2jIQzjqApy\", permissions: {\n publish: { allow: [\"$JS.API.>\", \"$JS.ACK.CONTROL.controller.>\", \"$JS.ACK.EVENTS.controller.>\", \"_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\"] }\n subscribe: { allow: [\"$JS.API.>\", \"_DELIVER.controller\", \"_DELIVER.controller.>\", \"_INBOX.controller.>\", \"mesh.control.>\", \"mesh.mod.mesh-catalog.event.catching-up\", \"mesh.mod.mesh-catalog.event.upgraded\", \"mesh.mod.gitea.event.pull.merged\", \"mesh.seat.mesh-build-machine.event.built\"] }\n allow_responses: { max: 1, ttl: \"1m\" }\n } }\n ]\n }\n}\n"
|
||||
},
|
||||
{
|
||||
"id": "broker",
|
||||
|
||||
@@ -2,12 +2,14 @@ module github.com/novox/mesh-host
|
||||
|
||||
go 1.26.0
|
||||
|
||||
require (
|
||||
github.com/nats-io/nats.go v1.54.0
|
||||
golang.org/x/crypto v0.57.0
|
||||
)
|
||||
|
||||
require (
|
||||
github.com/klauspost/compress v1.20.0 // indirect
|
||||
github.com/nats-io/nats.go v1.54.0 // indirect
|
||||
github.com/nats-io/nkeys v0.4.16 // indirect
|
||||
github.com/nats-io/nuid v1.0.1 // indirect
|
||||
github.com/rabbitmq/amqp091-go v1.14.0 // indirect
|
||||
golang.org/x/crypto v0.57.0 // indirect
|
||||
golang.org/x/sys v0.48.0 // indirect
|
||||
)
|
||||
|
||||
@@ -6,13 +6,7 @@ github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg=
|
||||
github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs=
|
||||
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
|
||||
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
|
||||
golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
|
||||
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
|
||||
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
|
||||
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
|
||||
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
|
||||
|
||||
+68
-1
@@ -320,6 +320,9 @@ func ApplyKeeping(
|
||||
// is a different thing — one is "this machine could not do it", the other is "this was never
|
||||
// a declaration", and they are fixed in different places.
|
||||
var failures []*Error
|
||||
// Modules whose own step did not complete. What follows *within such a module* is not attempted;
|
||||
// the rest of the machine is (novox/hq ADR 0136).
|
||||
gated := map[string]bool{}
|
||||
for i, resource := range ordered {
|
||||
if !orphansRemoved && i == guardFirst {
|
||||
// **Only a guard that is up may let the filter go.** Removing the derived filter's
|
||||
@@ -337,6 +340,17 @@ func ApplyKeeping(
|
||||
return report, known, err
|
||||
}
|
||||
}
|
||||
// **A module whose step did not complete is skipped from there on** (novox/hq ADR 0136).
|
||||
// Reported rather than passed over in silence: "not attempted" and "nothing to do" are
|
||||
// different answers, and only one of them is somebody's to fix.
|
||||
if module, ours := moduleOf(resource.Identity()); ours && gated[module] {
|
||||
report.Outcomes = append(report.Outcomes, Outcome{
|
||||
ID: resource.Identity(), Type: string(resource.Kind()), Target: resource.Target(),
|
||||
Action: "skipped", Detail: "a step this module declares did not complete",
|
||||
})
|
||||
continue
|
||||
}
|
||||
|
||||
// **On an adopted node, what is found is kept until its module is taken** (novox/hq ADR
|
||||
// 0100, ADR 0103). Before anything is applied: whatever of a module not yet taken is
|
||||
// present with no record of this host making it — or would reach what is — is held as it
|
||||
@@ -465,9 +479,26 @@ func ApplyKeeping(
|
||||
gates = true
|
||||
}
|
||||
if gates {
|
||||
failed.Gated = true
|
||||
// **A module's step gates that module, not the machine** (novox/hq ADR 0136).
|
||||
//
|
||||
// Stopping the whole apply is what this loop's own comment above calls holding a
|
||||
// machine hostage, and it was already rejected for every other shape (04-ISSUES/011).
|
||||
// A step exists to make something true before the next thing in *its module* needs it
|
||||
// — a store seeded before the broker starts, a schema prepared before the version
|
||||
// that needs it runs — so that is exactly how far the gate reaches. Everything else
|
||||
// on the machine is independent state and is attempted.
|
||||
//
|
||||
// An action still gates the machine: the bootstrap is a row of them, each making the
|
||||
// next possible, and they belong to no module.
|
||||
if module, ours := moduleOf(resource.Identity()); ours {
|
||||
gated[module] = true
|
||||
log(fmt.Sprintf(" gated %s.*: a step it declares did not complete, so the rest "+
|
||||
"of it was not attempted", module))
|
||||
continue
|
||||
}
|
||||
failed.Done = report
|
||||
failed.Others = len(failures) - 1
|
||||
failed.Gated = true
|
||||
return report, known, failed
|
||||
}
|
||||
continue
|
||||
@@ -1470,6 +1501,20 @@ func containerSpecReading(r *declaration.Container, declares, reads map[string]s
|
||||
for _, a := range r.Args {
|
||||
b.WriteString("arg " + a + "\n")
|
||||
}
|
||||
// **The mesh's names are part of what a container is** (novox/hq 04-ISSUES/135). A container
|
||||
// resolves every other machine and every public name through the entries the mesh gives it at
|
||||
// creation, and nothing re-reads them afterwards — so a container left alone when the roster
|
||||
// moved is one that cannot reach anything by name, for ever, while every check reports it
|
||||
// running. That is exactly what happened when this mesh's overlay range changed: one container
|
||||
// whose image and files never changed kept an address five days out of date and restarted
|
||||
// 2286 times against a database it could no longer find.
|
||||
//
|
||||
// Sorted, so the digest does not move for a reordering nobody made.
|
||||
hosts := append([]string(nil), r.Hosts...)
|
||||
sort.Strings(hosts)
|
||||
for _, h := range hosts {
|
||||
b.WriteString("host " + h + "\n")
|
||||
}
|
||||
// The resolver and address are part of what was declared: a container whose dns or ip moved
|
||||
// is a different container, or the fields could never reach one that already ran — which is
|
||||
// exactly how their first deployment silently changed nothing.
|
||||
@@ -2247,3 +2292,25 @@ func meshMadeUnits(known store.State) map[string]bool {
|
||||
}
|
||||
return made
|
||||
}
|
||||
|
||||
// moduleOf is the module a declared resource belongs to.
|
||||
//
|
||||
// The mesh composes a module's resource ids as `<module>.<its own id>`, and **a module's name may
|
||||
// contain a dot** — `novox.be` is one on this mesh — while a resource's own id never does. So the
|
||||
// owner is everything before the *last* dot; reading to the first one would make `novox.be.server`
|
||||
// belong to a module called "novox", and a gate would then skip whatever else happened to start that
|
||||
// way.
|
||||
//
|
||||
// False for what the mesh declares in its own right: the foundation's resources carry no dot at all,
|
||||
// and the adoption's are named for the mesh rather than for a module. Both belong to no module, and
|
||||
// their gate is therefore the machine's.
|
||||
func moduleOf(identity string) (string, bool) {
|
||||
if strings.HasPrefix(identity, declaration.AdoptionPrefix) {
|
||||
return "", false
|
||||
}
|
||||
at := strings.LastIndex(identity, ".")
|
||||
if at <= 0 {
|
||||
return "", false
|
||||
}
|
||||
return identity[:at], true
|
||||
}
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
package apply
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-host/internal/declaration"
|
||||
)
|
||||
|
||||
// **A container's mesh names are part of what it is** (novox/hq 04-ISSUES/135).
|
||||
//
|
||||
// A container resolves every machine and every public name through the entries it was given when it
|
||||
// was created, and nothing re-reads them. So a container the host leaves alone because nothing else
|
||||
// about it changed is a container that cannot reach anything by name — for ever, while every check
|
||||
// reports it running. That is what happened when this mesh's overlay range moved: one container kept
|
||||
// an address five days out of date and restarted 2286 times against a database it could no longer
|
||||
// find, and the host compared everything about it except that.
|
||||
func TestAContainersMeshNamesAreComparedLikeTheRestOfIt(t *testing.T) {
|
||||
was := &declaration.Container{
|
||||
Name: "umami", Image: "ghcr.io/example/umami@sha256:" + zeros(64),
|
||||
Hosts: []string{"novox.internal:10.42.0.1", "umami.novox.be:10.42.0.1"},
|
||||
}
|
||||
moved := &declaration.Container{
|
||||
Name: was.Name, Image: was.Image,
|
||||
Hosts: []string{"novox.internal:10.10.0.1", "umami.novox.be:10.10.0.1"},
|
||||
}
|
||||
if containerSpecReading(was, nil, nil) == containerSpecReading(moved, nil, nil) {
|
||||
t.Fatal("a container whose mesh names moved compares equal, so it is never recreated")
|
||||
}
|
||||
|
||||
// And the order they arrive in is not a change: the digest must not move for a reordering
|
||||
// nobody made.
|
||||
reordered := &declaration.Container{
|
||||
Name: moved.Name, Image: moved.Image,
|
||||
Hosts: []string{moved.Hosts[1], moved.Hosts[0]},
|
||||
}
|
||||
if containerSpecReading(moved, nil, nil) != containerSpecReading(reordered, nil, nil) {
|
||||
t.Fatal("the same names in another order read as a different container")
|
||||
}
|
||||
|
||||
// A container the mesh gives no names is unaffected, so nothing is recreated for a field it
|
||||
// does not set.
|
||||
plain := &declaration.Container{Name: "plex", Image: was.Image}
|
||||
if containerSpecReading(plain, nil, nil) == containerSpecReading(was, nil, nil) {
|
||||
return // different for other reasons, which is fine
|
||||
}
|
||||
}
|
||||
|
||||
func zeros(n int) string {
|
||||
out := make([]byte, n)
|
||||
for i := range out {
|
||||
out[i] = '0'
|
||||
}
|
||||
return string(out)
|
||||
}
|
||||
@@ -316,3 +316,85 @@ func TestAContainerNamingARunOnceStepIsRecreatedWhenItRan(t *testing.T) {
|
||||
t.Errorf("the recreation did not name the step as its reason: %+v", server)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAFailedStepGatesItsModuleAndNotTheMachine(t *testing.T) {
|
||||
// **The blast radius of a step is its module** (novox/hq ADR 0136). A step exists to make
|
||||
// something true before the next thing in its own module needs it — a store seeded before the
|
||||
// broker starts, a schema prepared before the version that needs it runs. Stopping the whole
|
||||
// apply is what this host's own loop calls holding a machine hostage, and it was already
|
||||
// rejected for every other shape (04-ISSUES/011): a module whose database is briefly
|
||||
// unreachable must not stop every module declared after it.
|
||||
var startedNames []string
|
||||
run := func(ctx context.Context, name string, args ...string) (string, error) {
|
||||
switch args[0] {
|
||||
case "info":
|
||||
return "27.0\n", nil
|
||||
case "container":
|
||||
return "false\t\n", errors.New("no such container")
|
||||
case "run":
|
||||
startedNames = append(startedNames, nameOf(args))
|
||||
if nameOf(args) == "catalogue-prepare" {
|
||||
return "", errors.New("exit status 1") // the schema could not be reached
|
||||
}
|
||||
return "deadbeef\n", nil
|
||||
case "rm":
|
||||
return "", nil
|
||||
}
|
||||
return "", nil
|
||||
}
|
||||
d := parseTrusted(t, `{"declaration":1,"resources":[
|
||||
{"id":"mesh-catalog.runtime-prepare","type":"container","name":"catalogue-prepare","image":"`+pinned+`","run-once":true},
|
||||
{"id":"mesh-catalog.runtime","type":"container","name":"catalogue","image":"`+pinned+`"},
|
||||
{"id":"gitea.server","type":"container","name":"forge","image":"`+pinned+`"}
|
||||
]}`)
|
||||
|
||||
report, _, err := Apply(context.Background(), archHost(t), d, store.State{}, store.OriginCarried, run, nil, nil)
|
||||
if err == nil {
|
||||
t.Fatal("a failed step was not reported as a failure")
|
||||
}
|
||||
started := map[string]bool{}
|
||||
for _, n := range startedNames {
|
||||
started[n] = true
|
||||
}
|
||||
if started["catalogue"] {
|
||||
t.Error("the module's own workload ran although its step did not complete")
|
||||
}
|
||||
if !started["forge"] {
|
||||
t.Error("another module was not attempted, so one module's step held the machine hostage")
|
||||
}
|
||||
// And the machine's own account says which was not attempted, rather than leaving it to be
|
||||
// inferred from silence.
|
||||
var skipped string
|
||||
for _, o := range report.Outcomes {
|
||||
if o.Action == "skipped" {
|
||||
skipped = o.ID
|
||||
}
|
||||
}
|
||||
if skipped != "mesh-catalog.runtime" {
|
||||
t.Errorf("the report does not say what was not attempted: %q", skipped)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAResourcesOwnerIsReadToTheLastDot(t *testing.T) {
|
||||
// **A module's name may contain a dot.** `novox.be` is one on this mesh, so reading a resource's
|
||||
// owner to the first dot would make its resources belong to something called "novox" — and a gate
|
||||
// would skip whatever else happened to start that way. A resource's own id never contains one,
|
||||
// which is what makes the last dot the boundary.
|
||||
for identity, want := range map[string]string{
|
||||
"novox.be.server": "novox.be",
|
||||
"mesh-catalog.runtime-prepare": "mesh-catalog",
|
||||
"gitea.admin-bootstrap": "gitea",
|
||||
} {
|
||||
got, ours := moduleOf(identity)
|
||||
if !ours || got != want {
|
||||
t.Errorf("%q belongs to %q (%v), want %q", identity, got, ours, want)
|
||||
}
|
||||
}
|
||||
// What the mesh declares in its own right belongs to no module: the foundation's resources carry
|
||||
// no dot, and the adoption's are the mesh's.
|
||||
for _, identity := range []string{"container-runtime", "store-ready", "adoption.guard", ".server"} {
|
||||
if _, ours := moduleOf(identity); ours {
|
||||
t.Errorf("%q was read as a module's", identity)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
@@ -16,9 +17,8 @@ import (
|
||||
// Approach is how a node reaches a mesh it does not yet belong to.
|
||||
//
|
||||
// The three things a token carries about where to go, and nothing about who is asking: the address,
|
||||
// the certificate that address must present, and which bus is at the other end. **Every token names
|
||||
// the bus the mesh runs on today until the rollout** (novox/hq ADR 0116 step 5), so an empty
|
||||
// Transport is the ordinary case rather than something missing.
|
||||
// the certificate that address must present, and which bus is at the other end — the mesh's own
|
||||
// (hearing.go), and a token naming any other is refused before anything is sent.
|
||||
type Approach struct {
|
||||
Address string
|
||||
Fingerprint string
|
||||
@@ -46,10 +46,8 @@ type Asking interface {
|
||||
func Present(ctx context.Context, to Approach, node, secret string,
|
||||
timeout time.Duration) (Asking, error) {
|
||||
|
||||
switch to.Transport {
|
||||
case OnNATS:
|
||||
return presentNats(ctx, to, node, secret, timeout)
|
||||
default:
|
||||
return presentCurrent(ctx, to, node, secret, timeout)
|
||||
if to.Transport != OnNATS {
|
||||
return nil, fmt.Errorf("this token is for the %q bus, and the mesh's bus is %s", to.Transport, OnNATS)
|
||||
}
|
||||
return presentNats(ctx, to, node, secret, timeout)
|
||||
}
|
||||
|
||||
@@ -1,129 +0,0 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The enrolment conversation on the bus the mesh runs on today.
|
||||
//
|
||||
// Moved out of Enrol rather than changed. The reply travels on this node's own queue and is picked
|
||||
// out by the correlation id the request carried, which is what this transport's reply field means
|
||||
// and has always meant.
|
||||
|
||||
type currentAsking struct {
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
queue string
|
||||
node string
|
||||
replies <-chan amqp.Delivery
|
||||
closed chan *amqp.Error
|
||||
}
|
||||
|
||||
func presentCurrent(_ context.Context, to Approach, node, secret string,
|
||||
timeout time.Duration) (Asking, error) {
|
||||
|
||||
config, err := PinnedConfig(to.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// The account name is the node's, and the password is the token's secret. Escaped because a name
|
||||
// or secret containing a colon or an at-sign would otherwise change which host this connects to
|
||||
// — a credential silently redirecting a connection is the worst shape this could take.
|
||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
||||
url.QueryEscape(node), url.QueryEscape(secret), to.Address)
|
||||
|
||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
||||
TLSClientConfig: config,
|
||||
Dial: amqp.DefaultDial(timeout),
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return nil, err
|
||||
}
|
||||
// Not quoted back: the DSN carries the one-time secret.
|
||||
return nil, fmt.Errorf("cannot reach the broker at %s as %s: %w", to.Address, node, err)
|
||||
}
|
||||
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// This node's own queue, which its account is scoped to and nothing else may read.
|
||||
queue, err := channel.QueueDeclare(QueueFor(node), true, false, false, false, nil)
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot declare this node's queue %s: %w", QueueFor(node), err)
|
||||
}
|
||||
|
||||
replies, err := channel.Consume(queue.Name, "", true, false, false, false, nil)
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return ¤tAsking{
|
||||
conn: conn, channel: channel, queue: queue.Name, node: node, replies: replies,
|
||||
closed: conn.NotifyClose(make(chan *amqp.Error, 1)),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (a *currentAsking) Close() {
|
||||
if a.channel != nil {
|
||||
_ = a.channel.Close()
|
||||
}
|
||||
if a.conn != nil {
|
||||
_ = a.conn.Close()
|
||||
}
|
||||
}
|
||||
|
||||
func (a *currentAsking) Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) {
|
||||
correlation := fmt.Sprintf("%s-%d", a.node, time.Now().UnixNano())
|
||||
publish, cancel := context.WithTimeout(ctx, wait)
|
||||
defer cancel()
|
||||
if err := a.channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false,
|
||||
amqp.Publishing{
|
||||
ContentType: "application/json",
|
||||
CorrelationId: correlation,
|
||||
ReplyTo: a.queue,
|
||||
Body: request,
|
||||
}); err != nil {
|
||||
return nil, fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err)
|
||||
}
|
||||
|
||||
// Waited for rather than assumed. A published message that nothing answers means the control
|
||||
// plane is not running, and a node that carried on regardless would believe it had joined a mesh
|
||||
// that has never heard of it.
|
||||
deadline := time.NewTimer(wait)
|
||||
defer deadline.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return nil, ctx.Err()
|
||||
case reason := <-a.closed:
|
||||
return nil, fmt.Errorf("the broker closed the connection: %v", reason)
|
||||
case <-deadline.C:
|
||||
return nil, fmt.Errorf(
|
||||
"the broker accepted this node's connection and nothing answered within %s. The "+
|
||||
"mesh's broker is running and its control plane is not", wait)
|
||||
case delivery, ok := <-a.replies:
|
||||
if !ok {
|
||||
return nil, errors.New("the broker stopped delivering")
|
||||
}
|
||||
// Anything else on this queue is not the answer to this question.
|
||||
if delivery.CorrelationId != correlation {
|
||||
continue
|
||||
}
|
||||
return delivery.Body, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -6,7 +6,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// Bus is what a host needs of the mesh's bus, in the mesh's own words.
|
||||
@@ -34,22 +33,6 @@ type Bus interface {
|
||||
|
||||
// --- The bus the mesh runs on today -----------------------------------------------------------
|
||||
|
||||
// OverCurrent is the bus as a channel, until the rollout.
|
||||
type OverCurrent struct{ Channel *amqp.Channel }
|
||||
|
||||
func (b OverCurrent) Report(ctx context.Context, node string, body []byte) error {
|
||||
// Mandatory: an unroutable report comes back rather than disappearing.
|
||||
return b.Channel.PublishWithContext(ctx, Exchange, KeyReport, true, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body})
|
||||
}
|
||||
|
||||
func (b OverCurrent) Alive(ctx context.Context, node string, body []byte) error {
|
||||
return b.Channel.PublishWithContext(ctx, Exchange, KeyAlive, false, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body})
|
||||
}
|
||||
|
||||
// --- NATS ---------------------------------------------------------------------------------
|
||||
|
||||
// OverNATS is the bus as a connection. A report goes through JetStream because it must survive
|
||||
// the controller's store restarting; a heartbeat does not, because it must not.
|
||||
type OverNATS struct {
|
||||
|
||||
+16
-25
@@ -2,15 +2,15 @@ package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
// What a host hears, as the host's own words for it.
|
||||
//
|
||||
// The outbound half went behind `Bus` (bus.go) and a node's two statements stopped naming a
|
||||
// transport. This is the other half — dialling, and the declarations that arrive — and it is where
|
||||
// the transport reached furthest: the run loop selected on a channel of the client library's own
|
||||
// delivery type, so every part of holding a node in its mesh knew which bus it was on.
|
||||
// The outbound half is behind `Bus` (bus.go); this is the other half — dialling, and the
|
||||
// declarations that arrive. The run loop reads its own words for a declaration rather than the
|
||||
// client library's delivery type, so nothing past this file knows what carried it.
|
||||
//
|
||||
// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005). This is its own
|
||||
// interface over its own libraries, and it agrees with the controller only because a conformance
|
||||
@@ -18,8 +18,8 @@ import (
|
||||
|
||||
// Link is this node's live connection to its mesh: what it hears, and what it says.
|
||||
//
|
||||
// One interface rather than two, because **dialling is where the transport is chosen** and choosing
|
||||
// it twice is how one half of a node ends up on a different bus from the other.
|
||||
// One interface rather than two, because dialling once is what keeps both halves of a node on the
|
||||
// same connection.
|
||||
type Link interface {
|
||||
// Bus is what this node says: what it applied, and that it is here.
|
||||
Bus
|
||||
@@ -45,9 +45,8 @@ type Link interface {
|
||||
// because applying is reconciliation: it converges rather than repeating.
|
||||
//
|
||||
// There is one way of being done rather than two. A declaration set aside because a newer arrived
|
||||
// with it is settled exactly as an applied one is, on both buses, and the difference between them
|
||||
// is a fact the *report* carries — a second method here would be a distinction the transport does
|
||||
// not make.
|
||||
// with it is settled exactly as an applied one is, and the difference between them is a fact the
|
||||
// *report* carries — a second method here would be a distinction the bus does not make.
|
||||
type Declaration interface {
|
||||
// Body is the signed declaration as it arrived, bytes unchanged: a node verifies what it
|
||||
// received rather than what it re-encoded.
|
||||
@@ -62,23 +61,15 @@ type Declaration interface {
|
||||
// Named Open rather than Dial because Dial is this package's raw TLS dial, which the enrolment path
|
||||
// uses to see a certificate before it trusts anything.
|
||||
//
|
||||
// **Both transports ship and this is the one place that chooses** (novox/hq ADR 0116: nothing moves
|
||||
// a node's bus before step 5). Until then every membership names the bus the mesh runs on today,
|
||||
// and the rollout is this switch and the credential behind it — not a change anywhere in the loop
|
||||
// that reads from what comes back.
|
||||
// **The mesh has one bus** (novox/hq ADR 0131): the one the broker seat delivers. A membership
|
||||
// still records which transport it was minted for, so a host can say what it is dialling, and a
|
||||
// membership recorded for anything else is a membership this host cannot use.
|
||||
func Open(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
||||
switch m.Transport {
|
||||
case OnNATS:
|
||||
return dialNats(ctx, m, timeout)
|
||||
default:
|
||||
return dialCurrent(ctx, m, timeout)
|
||||
if m.Transport != OnNATS {
|
||||
return nil, fmt.Errorf("this membership is for %q, and the mesh's bus is %s", m.Transport, OnNATS)
|
||||
}
|
||||
return dialNats(ctx, m, timeout)
|
||||
}
|
||||
|
||||
// The buses a node can be on. Empty is the one the mesh runs on today, which is every node until
|
||||
// the rollout — so a membership recorded before any of this existed reads as correct rather than as
|
||||
// unset.
|
||||
const (
|
||||
OnCurrent = ""
|
||||
OnNATS = "nats"
|
||||
)
|
||||
// OnNATS is the bus a membership names: the mesh's own, and the only one.
|
||||
const OnNATS = "nats"
|
||||
|
||||
@@ -1,164 +0,0 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// The host's link on the bus the mesh runs on today.
|
||||
//
|
||||
// Moved out of the run loop rather than changed: the dial, the queue, the prefetch window and the
|
||||
// return handler are what they were, because the mesh is running on this and a bus nothing speaks
|
||||
// yet is no reason to alter the one every node is on.
|
||||
|
||||
// currentLink is this node's connection as a channel.
|
||||
type currentLink struct {
|
||||
conn *amqp.Connection
|
||||
channel *amqp.Channel
|
||||
arrived chan Declaration
|
||||
lost chan error
|
||||
}
|
||||
|
||||
func dialCurrent(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
||||
config, err := PinnedConfig(m.Fingerprint)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
dsn := fmt.Sprintf("amqps://%s:%s@%s/",
|
||||
url.QueryEscape(m.Node), url.QueryEscape(m.Password), m.Broker)
|
||||
conn, err := amqp.DialConfig(dsn, amqp.Config{
|
||||
TLSClientConfig: config,
|
||||
Dial: amqp.DefaultDial(timeout),
|
||||
// Kept short so a node that has silently lost its route notices, rather than holding a
|
||||
// connection the broker forgot about and believing it is still in the mesh.
|
||||
Heartbeat: 10 * time.Second,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ErrWrongCertificate) {
|
||||
return nil, err
|
||||
}
|
||||
return nil, fmt.Errorf("cannot reach the broker at %s: %w", m.Broker, err)
|
||||
}
|
||||
|
||||
channel, err := conn.Channel()
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
queue := QueueFor(m.Node)
|
||||
if _, err := channel.QueueDeclare(queue, true, false, false, false, nil); err != nil {
|
||||
conn.Close()
|
||||
return nil, fmt.Errorf("cannot declare this node's queue %s: %w", queue, err)
|
||||
}
|
||||
|
||||
// Applying is one at a time — two at once would race on the same filesystem — but SEEING is
|
||||
// not: with a prefetch of one the host could never know that a newer declaration was already
|
||||
// waiting, and so applied every one of a backlog in turn, at the better part of a minute each,
|
||||
// becoming things nobody wanted any more (novox/hq issue 031). A window of unacknowledged
|
||||
// deliveries lets it drain to the newest; each declaration still survives a restart on the
|
||||
// broker until it is acknowledged, which happens only after it is applied or set aside.
|
||||
if err := channel.Qos(drainDepth, 0, false); err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
deliveries, err := channel.ConsumeWithContext(ctx, queue, "", false, false, false, false, nil)
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
l := ¤tLink{
|
||||
conn: conn, channel: channel,
|
||||
arrived: make(chan Declaration, drainDepth),
|
||||
lost: make(chan error, 1),
|
||||
}
|
||||
|
||||
// Published mandatory, so the broker hands back anything it cannot route rather than dropping
|
||||
// it. Without this a report goes to an exchange with no matching binding, the publisher is told
|
||||
// nothing, and the mesh believes this node never answered while the node believes it did —
|
||||
// which is what happened when `report` was left unbound on the other side.
|
||||
returned := channel.NotifyReturn(make(chan amqp.Return, 4))
|
||||
go func() {
|
||||
for r := range returned {
|
||||
select {
|
||||
case l.lost <- fmt.Errorf("the broker could not route this node's %s: %s (%d %s)",
|
||||
r.RoutingKey, r.Exchange, r.ReplyCode, r.ReplyText):
|
||||
default:
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
go func() {
|
||||
select {
|
||||
case reason := <-closed:
|
||||
select {
|
||||
case l.lost <- fmt.Errorf("the link closed: %v", reason):
|
||||
default:
|
||||
}
|
||||
case <-ctx.Done():
|
||||
}
|
||||
}()
|
||||
|
||||
// One goroutine turning the library's deliveries into the mesh's words, so the run loop selects
|
||||
// on one kind of thing whichever bus it is on.
|
||||
go func() {
|
||||
defer close(l.arrived)
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case delivery, ok := <-deliveries:
|
||||
if !ok {
|
||||
select {
|
||||
case l.lost <- errors.New("the broker stopped delivering"):
|
||||
default:
|
||||
}
|
||||
return
|
||||
}
|
||||
select {
|
||||
case l.arrived <- currentDeclaration{delivery}:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return l, nil
|
||||
}
|
||||
|
||||
func (l *currentLink) Declarations() <-chan Declaration { return l.arrived }
|
||||
func (l *currentLink) Lost() <-chan error { return l.lost }
|
||||
|
||||
func (l *currentLink) Close() {
|
||||
if l.channel != nil {
|
||||
_ = l.channel.Close()
|
||||
}
|
||||
if l.conn != nil {
|
||||
_ = l.conn.Close()
|
||||
}
|
||||
}
|
||||
|
||||
// Report and Alive are the outbound half, over the channel this link holds.
|
||||
func (l *currentLink) Report(ctx context.Context, node string, body []byte) error {
|
||||
return OverCurrent{Channel: l.channel}.Report(ctx, node, body)
|
||||
}
|
||||
|
||||
func (l *currentLink) Alive(ctx context.Context, node string, body []byte) error {
|
||||
return OverCurrent{Channel: l.channel}.Alive(ctx, node, body)
|
||||
}
|
||||
|
||||
// currentDeclaration is one delivery from the bus the mesh has.
|
||||
type currentDeclaration struct{ delivery amqp.Delivery }
|
||||
|
||||
func (d currentDeclaration) Body() []byte { return d.delivery.Body }
|
||||
func (d currentDeclaration) Handled() error { return d.delivery.Ack(false) }
|
||||
@@ -74,6 +74,9 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e
|
||||
// the server refused the bare name the first time a machine dialled it: "authentication
|
||||
// error - User". The same string the mesh composed into the user list, or nothing connects.
|
||||
nats.UserInfo("node."+m.Node, m.Password),
|
||||
// Replies to what this client asks the server arrive on its inbox, and the mesh grants a
|
||||
// machine exactly its own: the same prefix the user list was composed with.
|
||||
nats.CustomInboxPrefix("_INBOX.node." + m.Node),
|
||||
nats.Name("mesh-host/" + m.Node),
|
||||
nats.Timeout(timeout),
|
||||
// A node that has silently lost its route notices, rather than holding a connection the
|
||||
|
||||
@@ -110,6 +110,20 @@ type Report struct {
|
||||
// and the mesh's up in its place, and where the found configuration's original was kept.
|
||||
Tunnel *CarriedTunnel `json:"tunnel,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 a converged
|
||||
// node's filter is written around it.
|
||||
//
|
||||
// **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 this machine has no route off itself. The mesh then composes no filter for it and
|
||||
// leaves the one it has, rather than writing a rule around a link with no name — a rule set
|
||||
// that does not load is a machine filtering nothing while its unit reports success.
|
||||
Outward []string `json:"outward,omitempty"`
|
||||
|
||||
// Rekey is this node taking a found tunnel's key as its overlay key after enrolment (novox/hq
|
||||
// ADR 0105). Not an account of the machine: a report carrying one says nothing else.
|
||||
Rekey *Rekey `json:"rekey,omitempty"`
|
||||
|
||||
@@ -9,7 +9,7 @@ import (
|
||||
|
||||
// verified runs what Run does to a delivery body, without a broker: unmarshal, check the
|
||||
// signature, and only then apply. Isolating it keeps this test about the check rather than about
|
||||
// AMQP, which is tested against a real broker in the lab.
|
||||
// the bus, which is tested against a real one in the lab.
|
||||
func verified(t *testing.T, signer ed25519.PublicKey, body []byte) (Report, bool) {
|
||||
t.Helper()
|
||||
applied := false
|
||||
|
||||
@@ -30,9 +30,8 @@ type Membership struct {
|
||||
Fingerprint string
|
||||
Password string
|
||||
Signer ed25519.PublicKey
|
||||
// Transport is which bus this node speaks (hearing.go). Empty is the one the mesh runs on
|
||||
// today, which is every node until the rollout — so a membership recorded before any of this
|
||||
// existed reads as correct rather than as unset.
|
||||
// Transport is which bus this membership was minted for (hearing.go): the mesh's own, and a
|
||||
// membership that names another is one this host cannot dial with.
|
||||
Transport string
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,147 @@
|
||||
// Package outward reads which of this machine's links face outside it (novox/hq ADR 0140).
|
||||
//
|
||||
// The filter the mesh derives constrains traffic arriving from outside the machine and says nothing
|
||||
// about traffic that did not. To write that rule the mesh has to know which links "outside" arrives
|
||||
// on, and that is a thing only the machine can say — so it says it, once per report, the way it
|
||||
// already reports the kind of firewall it found and the tunnel it carried.
|
||||
//
|
||||
// **It replaces a list of addresses.** The filter used to allow the machine's own containers back
|
||||
// through by naming the address ranges they sit on: two ranges fixed in the control plane's source
|
||||
// and the rest typed by an operator. A range describes one machine and goes stale silently
|
||||
// (novox/hq 04-ISSUES/137 and /141). A link that carries the default route is a fact the machine
|
||||
// reads afresh every time, and it does not change when a module is added or removed.
|
||||
//
|
||||
// It reads the kernel's routing tables directly rather than asking a program. A module naming a
|
||||
// program the machine does not have is how the mesh already reported success while doing nothing
|
||||
// (novox/hq 04-ISSUES/136), and every machine has /proc.
|
||||
package outward
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// ProcNet is where the kernel publishes its routing tables. A parameter so a test can hold a
|
||||
// routing table without one.
|
||||
const ProcNet = "/proc/net"
|
||||
|
||||
// Links are the interfaces carrying a default route, for both address families, sorted and without
|
||||
// repeats.
|
||||
//
|
||||
// A machine may have more than one: a laptop with a cable and a radio has two, and both face
|
||||
// outside. A machine with none — no route off itself — returns nothing, and the mesh refuses to
|
||||
// compose a filter for it rather than writing a rule around a link with no name, which would be a
|
||||
// rule set that does not load and a machine filtering nothing while its unit reports success.
|
||||
func Links(procNet string) ([]string, error) {
|
||||
if procNet == "" {
|
||||
procNet = ProcNet
|
||||
}
|
||||
seen := map[string]bool{}
|
||||
|
||||
four, err := defaultsV4(filepath.Join(procNet, "route"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
six, err := defaultsV6(filepath.Join(procNet, "ipv6_route"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, name := range append(four, six...) {
|
||||
if name != "" && name != "lo" {
|
||||
seen[name] = true
|
||||
}
|
||||
}
|
||||
|
||||
out := make([]string, 0, len(seen))
|
||||
for name := range seen {
|
||||
out = append(out, name)
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// defaultsV4 reads /proc/net/route, whose columns are
|
||||
//
|
||||
// Iface Destination Gateway Flags RefCnt Use Metric Mask ...
|
||||
//
|
||||
// with addresses in hexadecimal. A default route is destination zero with mask zero — the mask
|
||||
// matters, because a route to the zero address with a real mask is not a default route.
|
||||
func defaultsV4(path string) ([]string, error) {
|
||||
lines, err := rows(path)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []string
|
||||
for _, fields := range lines {
|
||||
if len(fields) < 8 {
|
||||
continue
|
||||
}
|
||||
if isZeroHex(fields[1]) && isZeroHex(fields[7]) {
|
||||
out = append(out, fields[0])
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// defaultsV6 reads /proc/net/ipv6_route, whose columns are
|
||||
//
|
||||
// dest destprefix src srcprefix nexthop metric refcnt use flags iface
|
||||
//
|
||||
// A default route is the zero destination with a zero prefix length.
|
||||
func defaultsV6(path string) ([]string, error) {
|
||||
lines, err := rows(path)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []string
|
||||
for _, fields := range lines {
|
||||
if len(fields) < 10 {
|
||||
continue
|
||||
}
|
||||
if isZeroHex(fields[0]) && isZeroHex(fields[1]) {
|
||||
out = append(out, fields[9])
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// rows reads a routing table into fields per line, skipping a header and blank lines. A table that
|
||||
// is not there is not an error: a machine without the second address family has no file for it,
|
||||
// and that is not a machine that cannot be filtered.
|
||||
func rows(path string) ([][]string, error) {
|
||||
file, err := os.Open(path)
|
||||
if os.IsNotExist(err) {
|
||||
return nil, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot read the routing table at %s: %w", path, err)
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
var out [][]string
|
||||
scanner := bufio.NewScanner(file)
|
||||
for scanner.Scan() {
|
||||
line := strings.TrimSpace(scanner.Text())
|
||||
if line == "" || strings.HasPrefix(line, "Iface") {
|
||||
continue
|
||||
}
|
||||
out = append(out, strings.Fields(line))
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return nil, fmt.Errorf("cannot read the routing table at %s: %w", path, err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// isZeroHex is whether a hexadecimal field is all zeroes, whatever its width — the v4 table writes
|
||||
// eight digits and the v6 table thirty-two, and a prefix length is two.
|
||||
func isZeroHex(field string) bool {
|
||||
if field == "" {
|
||||
return false
|
||||
}
|
||||
return strings.Trim(strings.ToLower(field), "0") == ""
|
||||
}
|
||||
@@ -0,0 +1,120 @@
|
||||
package outward
|
||||
|
||||
import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// A routing table as the kernel writes it: a default route, a route to the zero address that is not
|
||||
// one, and a route on the loopback. Only the default route's link faces outside.
|
||||
const routeV4 = `Iface Destination Gateway Flags RefCnt Use Metric Mask MTU Window IRTT
|
||||
enp9s0 00000000 01FEA8C0 0003 0 0 100 00000000 0 0 0
|
||||
docker0 000011AC 00000000 0001 0 0 0 0000FFFF 0 0 0
|
||||
enp9s0 00000000 00000000 0001 0 0 100 00FFFFFF 0 0 0
|
||||
lo 00000000 00000000 0003 0 0 0 00000000 0 0 0
|
||||
`
|
||||
|
||||
const routeV6 = `00000000000000000000000000000000 00 00000000000000000000000000000000 00 fe800000000000000000000000000001 00000400 00000001 00000000 00000003 wlan0
|
||||
fd0000000000000000000000000000000 40 00000000000000000000000000000000 00 00000000000000000000000000000000 00000100 00000000 00000000 00000001 enp9s0
|
||||
`
|
||||
|
||||
func write(t *testing.T, dir, name, body string) {
|
||||
t.Helper()
|
||||
if err := os.WriteFile(filepath.Join(dir, name), []byte(body), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLinksAreTheOnesCarryingADefaultRoute(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
write(t, dir, "route", routeV4)
|
||||
write(t, dir, "ipv6_route", routeV6)
|
||||
|
||||
got, err := Links(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// The cable from the v4 table and the radio from the v6 one. Not docker0, whose route is not a
|
||||
// default; not the loopback, which faces nothing; and not the v6 route with a real prefix.
|
||||
want := []string{"enp9s0", "wlan0"}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("outward links are %v, want %v", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
// A route to the zero address with a real mask is not a default route. Trusting the destination
|
||||
// alone would name every link with such a route as facing outside, and a filter that treats an
|
||||
// internal bridge as outward constrains this machine's own guests — the fault ADR 0140 removes.
|
||||
func TestAZeroDestinationWithAMaskIsNotADefaultRoute(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
write(t, dir, "route", `Iface Destination Gateway Flags RefCnt Use Metric Mask MTU Window IRTT
|
||||
br-abc 00000000 00000000 0001 0 0 0 00FFFFFF 0 0 0
|
||||
`)
|
||||
got, err := Links(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got) != 0 {
|
||||
t.Fatalf("outward links are %v, want none", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A machine with no route off itself says so, rather than guessing. The mesh refuses to compose a
|
||||
// filter for it; a rule written around a link with no name does not load, and a rule set that does
|
||||
// not load is a machine filtering nothing while its unit reports success.
|
||||
func TestNoDefaultRouteIsNoLinks(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
write(t, dir, "route", "Iface\tDestination\tGateway \tFlags\tRefCnt\tUse\tMetric\tMask\t\tMTU\tWindow\tIRTT\n")
|
||||
got, err := Links(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got) != 0 {
|
||||
t.Fatalf("outward links are %v, want none", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A machine without the second address family has no file for it. That is not a machine that cannot
|
||||
// be filtered, so a missing table is read as no routes rather than as a failure.
|
||||
func TestAMissingTableIsNotAFailure(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
write(t, dir, "route", routeV4)
|
||||
got, err := Links(dir)
|
||||
if err != nil {
|
||||
t.Fatalf("a missing v6 table should not fail: %v", err)
|
||||
}
|
||||
if !reflect.DeepEqual(got, []string{"enp9s0"}) {
|
||||
t.Fatalf("outward links are %v, want [enp9s0]", got)
|
||||
}
|
||||
}
|
||||
|
||||
// The same link carrying a default route in both families is reported once.
|
||||
func TestALinkIsReportedOnce(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
write(t, dir, "route", routeV4)
|
||||
write(t, dir, "ipv6_route",
|
||||
"00000000000000000000000000000000 00 00000000000000000000000000000000 00 "+
|
||||
"fe800000000000000000000000000001 00000400 00000001 00000000 00000003 enp9s0\n")
|
||||
got, err := Links(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !reflect.DeepEqual(got, []string{"enp9s0"}) {
|
||||
t.Fatalf("outward links are %v, want [enp9s0]", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Against this machine's own routing table, so the parse is held to what the kernel actually writes
|
||||
// and not only to a fixture written to agree with it.
|
||||
func TestAgainstThisMachinesOwnTable(t *testing.T) {
|
||||
got, err := Links("")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got) == 0 {
|
||||
t.Skip("this machine has no default route")
|
||||
}
|
||||
t.Logf("this machine's outward links: %v", got)
|
||||
}
|
||||
Reference in New Issue
Block a user