From 2478b5f127959333dd97ba3b4e59c97a95c5a748 Mon Sep 17 00:00:00 2001 From: jochen Date: Mon, 28 Sep 2026 03:54:22 +0200 Subject: [PATCH] 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. --- cmd/mesh-host/main.go | 6 +- go.mod | 8 +- go.sum | 6 -- internal/link/asking.go | 14 ++- internal/link/asking_current.go | 129 ------------------------ internal/link/bus.go | 17 ---- internal/link/hearing.go | 41 +++----- internal/link/hearing_current.go | 164 ------------------------------- internal/link/messages_test.go | 2 +- internal/link/run.go | 5 +- 10 files changed, 31 insertions(+), 361 deletions(-) delete mode 100644 internal/link/asking_current.go delete mode 100644 internal/link/hearing_current.go diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 0afbc03..5b6ae72 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -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) } diff --git a/go.mod b/go.mod index 47f3aaf..2c02e75 100644 --- a/go.mod +++ b/go.mod @@ -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 ) diff --git a/go.sum b/go.sum index 726180b..65100b8 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/link/asking.go b/internal/link/asking.go index 529d9ae..2a53261 100644 --- a/internal/link/asking.go +++ b/internal/link/asking.go @@ -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) } diff --git a/internal/link/asking_current.go b/internal/link/asking_current.go deleted file mode 100644 index d1fb792..0000000 --- a/internal/link/asking_current.go +++ /dev/null @@ -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 - } - } -} diff --git a/internal/link/bus.go b/internal/link/bus.go index 7d0f9b7..540ff89 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -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 { diff --git a/internal/link/hearing.go b/internal/link/hearing.go index 303d7e7..173ea12 100644 --- a/internal/link/hearing.go +++ b/internal/link/hearing.go @@ -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" diff --git a/internal/link/hearing_current.go b/internal/link/hearing_current.go deleted file mode 100644 index 772e17e..0000000 --- a/internal/link/hearing_current.go +++ /dev/null @@ -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) } diff --git a/internal/link/messages_test.go b/internal/link/messages_test.go index 76a9a7e..a35d52a 100644 --- a/internal/link/messages_test.go +++ b/internal/link/messages_test.go @@ -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 diff --git a/internal/link/run.go b/internal/link/run.go index 0288a4f..47c5edf 100644 --- a/internal/link/run.go +++ b/internal/link/run.go @@ -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 }