Author SHA1 Message Date
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
mesh-admin c82e933268 Merge pull request 'Genesis: the first user list lets the controller hear the forge's merges' (#38) from fix/genesis-follows-the-forge into main 2026-09-28 01:13:18 +00:00
jschoubben e212da6bd2 Genesis: the first user list lets the controller hear the forge's merges 2026-09-28 03:13:16 +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
jschoubben 68a193d6f0 A host adopts a delivered membership at start, not only after a declaration
The rescue path: an operator writes the membership file by hand on a machine no
bus can reach — rotated while it still held the old password — and restarts the
host. Same file, same check as the declaration path.
2026-09-28 02:21:56 +02:00
jschoubben 4fba884d47 Merge pull request 'Genesis: the first user list lets the controller hear its consumers, and the account has JetStream' (#36) from fix/genesis-accounts-hear-and-have-jetstream into main 2026-09-27 23:47:41 +00:00
jschoubben 15f0dabf32 Genesis: the first user list lets the controller hear its consumers, and the account has JetStream
Two things the controller's own composition now derives and the installer's
carried list did not: the delivery subjects of the controller's consumers, and
JetStream enabled on the account — without which the first bound consumer is
refused. Found live on a mesh moved rather than raised; a fresh genesis would
have met both at first start.
2026-09-28 01:46:14 +02:00
jschoubben 4ef41ad5a0 Merge pull request 'A host logs in as node.<name> and hears answers on its own inbox' (#35) from fix/a-host-uses-its-own-inbox into main 2026-09-27 23:18:30 +00:00
jschoubben d749de989c A host hears the server's answers on its own inbox
The mesh grants a machine exactly its own inbox prefix; the client used the
default one and was refused a subscription to it.
2026-09-28 01:16:23 +02:00
jschoubben 14e64844df A host logs in to the new bus as node.<name>
The mesh names a machine's bus user node.<name>; the host dialled with the bare
name and every machine was refused the first time it reached the server:
"authentication error - User". Same string as the composed user list, or nothing
connects.
2026-09-28 01:13:52 +02:00
jschoubben d82fc9a121 Merge pull request 'A machine moves to the bus its declaration tells it to' (#34) from feat/a-machine-moves-to-the-bus-it-is-told into main 2026-09-27 22:09:08 +00:00
jschoubben 9fd3762e93 A machine moves to the bus its declaration tells it to
Until now a membership — bus address, fingerprint, password — was written once,
at enrolment, and nothing ever rewrote it, so a machine already enrolled could not
be moved to another bus at all (novox/hq design 28, task 5.2). The mesh now
delivers a membership for the new bus as a sealed file in the declaration, like
any secret; the host reads it after the declaration has applied, saves it as its
identity, and exits cleanly so the service manager restarts it dialling the bus it
names. The same path a machine takes after a reboot — so nothing new has to be
right for it to work. A membership carries its transport; empty means the bus the
mesh ran on before, so nothing written earlier reads as unset.
2026-09-28 00:08:44 +02:00
jschoubben 587b8aa220 Merge pull request 'A host reaches either bus, and a genesis template that raises the mesh on the new one' (#33) from feat/nats-genesis into main 2026-09-27 17:36:48 +00:00
13 changed files with 106 additions and 361 deletions
+63 -3
View File
@@ -869,6 +869,7 @@ func overlayCommand(ctx context.Context, opts options) error {
if err := link.Publish(ctx, link.Membership{ if err := link.Publish(ctx, link.Membership{
Node: mine.Node, Broker: mine.Membership.Broker, Fingerprint: mine.Membership.Fingerprint, Node: mine.Node, Broker: mine.Membership.Broker, Fingerprint: mine.Membership.Fingerprint,
Password: mine.Membership.Password, Signer: mine.Membership.Signer, Password: mine.Membership.Password, Signer: mine.Membership.Signer,
Transport: mine.Membership.Transport,
}, link.Report{Node: mine.Node, Rekey: &rekey}, opts.timeout); err != nil { }, link.Report{Node: mine.Node, Rekey: &rekey}, opts.timeout); err != nil {
return fmt.Errorf("the mesh could not be told; nothing was written here: %w", err) return fmt.Errorf("the mesh could not be told; nothing was written here: %w", err)
} }
@@ -961,9 +962,14 @@ func runLink(ctx context.Context, opts options) error {
return nil // asked to stop while waiting 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) } 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 // 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 // (novox/hq ADR 0053). It fires scheduled steps on their cadence, surviving across applies and
@@ -973,7 +979,12 @@ func runLink(ctx context.Context, opts options) error {
go sched.Run(ctx) go sched.Run(ctx)
applier := func(ctx context.Context, raw, signature []byte) link.Report { applier := func(ctx context.Context, raw, signature []byte) link.Report {
return applyAndKeep(ctx, opts, raw, &store.Declared{Declaration: raw, Signature: signature}, sched, say) report := applyAndKeep(ctx, opts, raw, &store.Declared{Declaration: raw, Signature: signature}, sched, say)
// **A declaration may carry this machine's membership for another bus.** It arrives as a
// sealed file like any secret, and is read after the rest has applied so the bus it names is
// standing before this machine leaves the one it is on (novox/hq design 28, task 5.2).
adoptDeliveredMembership(identity.Path(opts.state), &mine, say)
return report
} }
// Two things at once, and the second is what makes disconnection ordinary. The link brings // Two things at once, and the second is what makes disconnection ordinary. The link brings
@@ -1008,6 +1019,7 @@ func runLink(ctx context.Context, opts options) error {
Broker: mine.Membership.Broker, Broker: mine.Membership.Broker,
Fingerprint: mine.Membership.Fingerprint, Fingerprint: mine.Membership.Fingerprint,
Password: mine.Membership.Password, Password: mine.Membership.Password,
Transport: mine.Membership.Transport,
Signer: mine.Membership.Signer, Signer: mine.Membership.Signer,
}, applier, say, opts.timeout, rousedBySignal(ctx), outbox) }, applier, say, opts.timeout, rousedBySignal(ctx), outbox)
} }
@@ -1376,3 +1388,51 @@ func carriedPorts(state store.State) []int {
sort.Ints(out) sort.Ints(out)
return out return out
} }
// MembershipNextPath is where the mesh delivers this machine's membership for the bus it is moving
// to: a sealed file in a declaration, written by the host like any secret, read here after the
// declaration has applied.
const MembershipNextPath = "/var/lib/mesh/membership-next.json"
// adoptDeliveredMembership moves this machine to the bus a delivered membership names.
//
// **Saved, then restarted — not swapped in place.** The link holds one membership for the life of
// the process, and every reconnect path assumes the bus did not change under it; a process that
// found itself half on one bus and half on another would be a new kind of state nothing was written
// for. Exiting cleanly hands the machine to the service manager's restart, and the process that
// comes back reads the identity file the way it always has and dials the bus it names. That is the
// same path a machine takes after a reboot, which is why nothing new has to be right for it to work.
//
// A membership identical to the one held is nothing: the file stays and is read again next time.
func adoptDeliveredMembership(identityPath string, mine *identity.Identity, say func(string)) {
raw, err := os.ReadFile(MembershipNextPath)
if err != nil {
return
}
var next identity.Membership
if err := json.Unmarshal(raw, &next); err != nil {
say(fmt.Sprintf("a membership was delivered at %s and could not be read: %v", MembershipNextPath, err))
return
}
if next.Broker == "" || next.Password == "" || next.Fingerprint == "" {
say(fmt.Sprintf("a membership was delivered at %s with no broker, fingerprint or password; ignored", MembershipNextPath))
return
}
same := next.Broker == mine.Membership.Broker && next.Fingerprint == mine.Membership.Fingerprint &&
next.Password == mine.Membership.Password && next.Transport == mine.Membership.Transport
if same {
return
}
// The signer is the mesh's, not the bus's: a delivered membership that names none keeps the one
// this machine already trusts, because a change of bus is not a change of who signs declarations.
if len(next.Signer) == 0 {
next.Signer = mine.Membership.Signer
}
mine.Membership = next
if err := identity.Save(identityPath, *mine); err != nil {
say(fmt.Sprintf("a membership for another bus was delivered and could not be saved: %v", err))
return
}
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", "type": "file",
"path": "/var/lib/mesh-bus-conf/accounts.conf", "path": "/var/lib/mesh-bus-conf/accounts.conf",
"mode": "0600", "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.>\"] }\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", "id": "broker",
+5 -3
View File
@@ -2,12 +2,14 @@ module github.com/novox/mesh-host
go 1.26.0 go 1.26.0
require (
github.com/nats-io/nats.go v1.54.0
golang.org/x/crypto v0.57.0
)
require ( require (
github.com/klauspost/compress v1.20.0 // indirect 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/nkeys v0.4.16 // indirect
github.com/nats-io/nuid v1.0.1 // 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 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/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 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= 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 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA= 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 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og= golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
+5
View File
@@ -76,6 +76,11 @@ type Membership struct {
// Not the token's secret: that is spent, and a credential that lives for ever should not be // Not the token's secret: that is spent, and a credential that lives for ever should not be
// the same string as one that was meant to be used once. // the same string as one that was meant to be used once.
Password string `json:"password"` Password string `json:"password"`
// Transport is which bus this membership is for. Empty is the bus the mesh ran on before
// the move — so every membership written before this field existed reads as correct, not as
// unset — and "nats" is the one being moved to (novox/hq design 28, task 5.2).
Transport string `json:"transport,omitempty"`
} }
// Queue is where this node listens. Its account may read this and nothing else. // Queue is where this node listens. Its account may read this and nothing else.
+6 -8
View File
@@ -2,6 +2,7 @@ package link
import ( import (
"context" "context"
"fmt"
"time" "time"
) )
@@ -16,9 +17,8 @@ import (
// Approach is how a node reaches a mesh it does not yet belong to. // 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 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 certificate that address must present, and which bus is at the other end — the mesh's own
// the bus the mesh runs on today until the rollout** (novox/hq ADR 0116 step 5), so an empty // (hearing.go), and a token naming any other is refused before anything is sent.
// Transport is the ordinary case rather than something missing.
type Approach struct { type Approach struct {
Address string Address string
Fingerprint string Fingerprint string
@@ -46,10 +46,8 @@ type Asking interface {
func Present(ctx context.Context, to Approach, node, secret string, func Present(ctx context.Context, to Approach, node, secret string,
timeout time.Duration) (Asking, error) { timeout time.Duration) (Asking, error) {
switch to.Transport { if to.Transport != OnNATS {
case 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)
default:
return presentCurrent(ctx, to, node, secret, timeout)
} }
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" "time"
"github.com/nats-io/nats.go" "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. // 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 ----------------------------------------------------------- // --- 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 // 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. // the controller's store restarting; a heartbeat does not, because it must not.
type OverNATS struct { type OverNATS struct {
+16 -25
View File
@@ -2,15 +2,15 @@ package link
import ( import (
"context" "context"
"fmt"
"time" "time"
) )
// What a host hears, as the host's own words for it. // 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 // The outbound half is behind `Bus` (bus.go); this is the other half — dialling, and the
// transport. This is the other half — dialling, and the declarations that arrive — and it is where // declarations that arrive. The run loop reads its own words for a declaration rather than the
// the transport reached furthest: the run loop selected on a channel of the client library's own // client library's delivery type, so nothing past this file knows what carried it.
// delivery type, so every part of holding a node in its mesh knew which bus it was on.
// //
// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005). This is its own // **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 // 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. // 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 // One interface rather than two, because dialling once is what keeps both halves of a node on the
// it twice is how one half of a node ends up on a different bus from the other. // same connection.
type Link interface { type Link interface {
// Bus is what this node says: what it applied, and that it is here. // Bus is what this node says: what it applied, and that it is here.
Bus Bus
@@ -45,9 +45,8 @@ type Link interface {
// because applying is reconciliation: it converges rather than repeating. // 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 // 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 // with it is settled exactly as an applied one is, and the difference between them is a fact the
// is a fact the *report* carries — a second method here would be a distinction the transport does // *report* carries — a second method here would be a distinction the bus does not make.
// not make.
type Declaration interface { type Declaration interface {
// Body is the signed declaration as it arrived, bytes unchanged: a node verifies what it // Body is the signed declaration as it arrived, bytes unchanged: a node verifies what it
// received rather than what it re-encoded. // 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 // 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. // 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 // **The mesh has one bus** (novox/hq ADR 0131): the one the broker seat delivers. A membership
// a node's bus before step 5). Until then every membership names the bus the mesh runs on today, // still records which transport it was minted for, so a host can say what it is dialling, and a
// and the rollout is this switch and the credential behind it — not a change anywhere in the loop // membership recorded for anything else is a membership this host cannot use.
// that reads from what comes back.
func Open(ctx context.Context, m Membership, timeout time.Duration) (Link, error) { func Open(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
switch m.Transport { if m.Transport != OnNATS {
case 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)
default:
return dialCurrent(ctx, m, timeout)
} }
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 // OnNATS is the bus a membership names: the mesh's own, and the only one.
// the rollout — so a membership recorded before any of this existed reads as correct rather than as const OnNATS = "nats"
// unset.
const (
OnCurrent = ""
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) }
+7 -1
View File
@@ -70,7 +70,13 @@ func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, e
} }
opts := []nats.Option{ opts := []nats.Option{
nats.Secure(config), nats.Secure(config),
nats.UserInfo(m.Node, m.Password), // The mesh names a machine's bus user "node.<name>" (the controller's principal scheme), and
// 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.Name("mesh-host/" + m.Node),
nats.Timeout(timeout), nats.Timeout(timeout),
// A node that has silently lost its route notices, rather than holding a connection the // A node that has silently lost its route notices, rather than holding a connection the
+1 -1
View File
@@ -9,7 +9,7 @@ import (
// verified runs what Run does to a delivery body, without a broker: unmarshal, check the // 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 // 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) { func verified(t *testing.T, signer ed25519.PublicKey, body []byte) (Report, bool) {
t.Helper() t.Helper()
applied := false applied := false
+2 -3
View File
@@ -30,9 +30,8 @@ type Membership struct {
Fingerprint string Fingerprint string
Password string Password string
Signer ed25519.PublicKey Signer ed25519.PublicKey
// Transport is which bus this node speaks (hearing.go). Empty is the one the mesh runs on // Transport is which bus this membership was minted for (hearing.go): the mesh's own, and a
// today, which is every node until the rollout — so a membership recorded before any of this // membership that names another is one this host cannot dial with.
// existed reads as correct rather than as unset.
Transport string Transport string
} }