10 Commits
Author SHA1 Message Date
mesh-admin 99562947f3 Merge pull request 'Genesis lets the control plane state its own facts' (#43) from fix/genesis-lets-the-control-plane-state-its-facts into main 2026-09-28 14:07:22 +00:00
jschoubben c171c64d3a Genesis lets the control plane state its own facts
A mesh raised from nothing must be able to say what it applied and what a machine refused (novox/hq
ADR 0134), and the first user list is what permits it. Kept in step with the controller's own
composition by the test that reads this file.
2026-09-28 16:07:20 +02:00
mesh-admin 0dc5515099 Merge pull request 'A resource's owner is read to the last dot, because a module's name may contain one' (#42) from fix/a-resources-owner-is-read-to-the-last-dot into main 2026-09-28 13:45:01 +00:00
jschoubben 58c0e715c7 A resource's owner is read to the last dot, because a module's name may contain one
novox.be is a module on this mesh. Reading a resource's owner to the first dot made its resources
belong to something called "novox", and a step's gate would then skip whatever else happened to start
that way — silently. A resource's own id never contains a dot, which is what makes the last one the
boundary; the mesh's derived step was changed to add a hyphen rather than a dot for the same reason
(mesh-controller #128).
2026-09-28 15:44:59 +02:00
mesh-admin c5b229f95d Merge pull request 'A step gates its module, not the machine' (#41) from fix/a-step-gates-its-module-not-the-machine into main 2026-09-28 13:38:21 +00:00
jschoubben a9c49724b5 A step gates its module, not the machine
A run-once container that does not complete stops the rest of that module's resources; everything
else on the machine is attempted, as every other shape already is (novox/hq ADR 0136). An action
still gates the machine — genesis is a row of them and they belong to no module. What was not
attempted is reported as skipped, because that and 'nothing to do' are different answers.

Without this, ADR 0135's derived preparation would let one module's unreachable database hold a
machine hostage — the fault 04-ISSUES/011 removed for everything else, and the reason the catalogue
migrates itself at start.
2026-09-28 15:38:19 +02:00
mesh-admin 296ec5ece6 Merge pull request 'Genesis lets the controller ask any module's tool' (#40) from fix/genesis-lets-the-controller-ask into main 2026-09-28 02:21:24 +00:00
jschoubben 45b9a507a1 Genesis lets the controller ask any module's tool
The controller is the way in for tool calls (novox/hq ADR 0095); the first user list must say so,
or a mesh raised from nothing refuses its own first ask. Kept in step with the controller's
composition by the test that reads this file.
2026-09-28 04:21:22 +02:00
mesh-admin 5c410aaf54 Merge pull request 'One bus: the AMQP transport is gone from the host' (#39) from feat/one-bus into main 2026-09-28 01:54:23 +00:00
jschoubben 2478b5f127 One bus: the AMQP transport is gone from the host
The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The host's old
dialling and enrolment paths are deleted with the switch that chose between them; a membership or
a token naming another bus is refused before anything is sent, rather than dialled on a transport
that no longer exists.
2026-09-28 03:54:22 +02:00
13 changed files with 168 additions and 363 deletions
+1 -5
View File
@@ -1433,10 +1433,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)
}
+1 -1
View File
@@ -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 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.node.>\", \"mesh.seat.mesh-build-machine.accept.>\"] }\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"
"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",
+5 -3
View File
@@ -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
View File
@@ -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=
+54 -1
View File
@@ -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
@@ -2247,3 +2278,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
}
+82
View File
@@ -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)
}
}
}
+6 -8
View File
@@ -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)
}
-129
View File
@@ -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 &currentAsking{
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
}
}
}
-17
View File
@@ -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
View File
@@ -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"
-164
View File
@@ -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 := &currentLink{
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) }
+1 -1
View File
@@ -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
+2 -3
View File
@@ -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
}