Author SHA1 Message Date
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
mesh-admin 673912ccdc Merge pull request 'A host adopts a delivered membership at start, not only after a declaration' (#37) from feat/a-host-adopts-a-membership-at-start into main 2026-09-28 00:40:27 +00:00
10 changed files with 31 additions and 361 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)
}
+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=
+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
}