diff --git a/cmd/mesh-host/main.go b/cmd/mesh-host/main.go index 3bacc00..1c83772 100644 --- a/cmd/mesh-host/main.go +++ b/cmd/mesh-host/main.go @@ -772,7 +772,12 @@ func enrol(ctx context.Context, opts options) error { // who knows its public key (novox/hq issue 083). proof := mine.Sign(link.EnrolProof(token.Secret, mine.Public, mine.Overlay.Public, sealing.Public, serving.Public)) - reply, err := link.Enrol(ctx, token.Broker, token.Fingerprint, *name, token.Secret, + // The token says where to go and which certificate that address must present. It says nothing + // about which bus is there, and does not need to: every token names the one the mesh runs on + // today until the rollout (novox/hq ADR 0116 step 5), and that is what an empty Transport is. + reply, err := link.Enrol(ctx, + link.Approach{Address: token.Broker, Fingerprint: token.Fingerprint}, + *name, token.Secret, mine.Public, mine.Overlay.Public, sealing.Public, serving.Public, reported, proof, found, opts.timeout) if err != nil { diff --git a/internal/link/asking.go b/internal/link/asking.go new file mode 100644 index 0000000..529d9ae --- /dev/null +++ b/internal/link/asking.go @@ -0,0 +1,55 @@ +package link + +import ( + "context" + "time" +) + +// The enrolment conversation, as the host's own words for it. +// +// **Its own seam rather than part of Link**, because almost nothing about it is the same. The +// credential is a one-time secret rather than this node's own; there is no declaration to hear; the +// whole exchange is a single question asked and possibly asked again. And the stakes differ: a node +// that fails here is not in the mesh at all, where a node that fails in Link has merely lost touch +// with one it belongs 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 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. +type Approach struct { + Address string + Fingerprint string + Transport string +} + +// Asking is one open enrolment conversation. +type Asking interface { + // Ask puts the request to the mesh and waits for one answer, or says why none came. + // + // Called again, with the same bytes, while the mesh says "try again": the keys this node + // generated are the ones it keeps, so the same request is the same enrolment and the mesh holds + // the token for it (novox/hq issue 083). + Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) + + // Close lets go of the connection made with the token. + Close() +} + +// Present opens an enrolment conversation with the mesh. +// +// The connection is made before anything is sent, and the certificate is checked while it is being +// made — so a node pointed at the wrong bus finds out before its token has left the machine (ADR +// 0004). +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) + } +} diff --git a/internal/link/asking_current.go b/internal/link/asking_current.go new file mode 100644 index 0000000..d1fb792 --- /dev/null +++ b/internal/link/asking_current.go @@ -0,0 +1,129 @@ +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/asking_nats.go b/internal/link/asking_nats.go new file mode 100644 index 0000000..b8df9e6 --- /dev/null +++ b/internal/link/asking_nats.go @@ -0,0 +1,166 @@ +package link + +import ( + "context" + "crypto/rand" + "encoding/hex" + "errors" + "fmt" + "time" + + "github.com/nats-io/nats.go" +) + +// The enrolment conversation on the bus being built. +// +// **The reply address is the whole of what changes**, and it changes for a reason the transport +// forces rather than a preference. Core NATS request/reply puts the caller's inbox in the message's +// reply field and a plain responder answers it — but this request goes into a stream, and a message +// a JetStream consumer delivers has had that field claimed for the consumer's own ack address. So by +// the time the controller reads the request, the transport's reply field names where the +// *controller* must acknowledge. Verified against a running server (design 25 §2). +// +// The address therefore travels as a field of the request, and this subscribes it before publishing: +// a node that published first could miss an answer to a question nobody was listening for. + +// enrolInbox is where a node enrolling waits. +// +// Under `_INBOX.enrol..`, which is exactly what its enrolment user may subscribe and no +// wider — so an answer sealed to one machine cannot be read by another enrolling beside it. The +// random tail is this attempt's own: a reply left over from an attempt that timed out is not the +// answer to this question, which is what the correlation id does on the other transport. +func enrolInbox(node string) (string, error) { + tail := make([]byte, 8) + if _, err := rand.Read(tail); err != nil { + return "", fmt.Errorf("cannot make a reply address: %w", err) + } + return "_INBOX.enrol." + node + "." + hex.EncodeToString(tail), nil +} + +type natsAsking struct { + conn *nats.Conn + js nats.JetStreamContext + inbox string + answers *nats.Subscription + lost chan error +} + +func presentNats(_ context.Context, to Approach, node, secret string, + timeout time.Duration) (Asking, error) { + + config, err := PinnedConfig(to.Fingerprint) + if err != nil { + return nil, err + } + + inbox, err := enrolInbox(node) + if err != nil { + return nil, err + } + + lost := make(chan error, 1) + // The user is this token's own — `enrol.`, which may publish the enrolment subject and + // subscribe its own inbox and nothing else (design 25 §6). The secret is its password, the same + // string the request claims, so the server proves somebody holds the token and the request + // proves the same thing to the controller without it having to ask the server who connected. + conn, err := nats.Connect(natsURL(to.Address), + nats.Secure(config), + nats.UserInfo("enrol."+node, secret), + nats.Name("mesh-host/enrol/"+node), + nats.Timeout(timeout), + nats.NoReconnect(), + nats.DisconnectErrHandler(func(_ *nats.Conn, err error) { + select { + case lost <- fmt.Errorf("the bus closed the connection: %w", err): + default: + } + }), + ) + if err != nil { + if errors.Is(err, ErrWrongCertificate) { + return nil, err + } + // Not quoted back with the credential: the secret is one-time and still a secret. + return nil, fmt.Errorf("cannot reach the bus at %s as %s: %w", to.Address, node, err) + } + js, err := conn.JetStream() + if err != nil { + conn.Close() + return nil, fmt.Errorf("the bus at %s has no JetStream: %w", to.Address, err) + } + + // Subscribed before anything is published, so an answer cannot arrive before there is anywhere + // for it to land. + answers, err := conn.SubscribeSync(inbox) + if err != nil { + conn.Close() + return nil, fmt.Errorf("this node cannot listen for the mesh's answer: %w", err) + } + if err := conn.Flush(); err != nil { + conn.Close() + return nil, fmt.Errorf("this node's reply address did not reach the bus: %w", err) + } + + return &natsAsking{conn: conn, js: js, inbox: inbox, answers: answers, lost: lost}, nil +} + +func (a *natsAsking) Close() { + if a.answers != nil { + _ = a.answers.Unsubscribe() + } + if a.conn != nil { + a.conn.Close() + } +} + +// Ask publishes the request with this attempt's reply address written into it, and waits there. +func (a *natsAsking) Ask(ctx context.Context, request []byte, wait time.Duration) ([]byte, error) { + addressed, err := withReplyTo(request, a.inbox) + if err != nil { + return nil, err + } + + publish, cancel := context.WithTimeout(ctx, wait) + defer cancel() + // Into the stream and awaited: an enrolment the bus never accepted must fail here rather than be + // assumed, because the node has nothing else to go on. + if _, err := a.js.Publish(EnrolSubject, addressed, nats.Context(publish)); err != nil { + return nil, fmt.Errorf("cannot ask the mesh to enrol this node: %w", err) + } + + // Waited for rather than assumed. A published message that nothing answers means the controller + // is not running, and a node that carried on regardless would believe it had joined a mesh that + // has never heard of it. + // + // **The wait may legitimately be several store-window cycles long**: the controller naks the + // request with a delay while its store is restarting, and the node is waiting on the other side + // of that — which is exactly the combination that would have delivered the answer to a caller + // who had given up, had the address travelled in the transport's field. + answered, cancelAnswer := context.WithTimeout(ctx, wait) + defer cancelAnswer() + for { + msg, err := a.answers.NextMsgWithContext(answered) + switch { + case err == nil: + return msg.Data, nil + case errors.Is(err, context.DeadlineExceeded): + select { + case reason := <-a.lost: + return nil, reason + default: + } + return nil, fmt.Errorf( + "the bus accepted this node's connection and nothing answered within %s. The mesh's "+ + "bus is running and its controller is not", wait) + case errors.Is(err, context.Canceled): + return nil, ctx.Err() + default: + select { + case reason := <-a.lost: + return nil, reason + default: + } + return nil, fmt.Errorf("waiting for the mesh's answer: %w", err) + } + } +} diff --git a/internal/link/asking_nats_test.go b/internal/link/asking_nats_test.go new file mode 100644 index 0000000..5d853b2 --- /dev/null +++ b/internal/link/asking_nats_test.go @@ -0,0 +1,135 @@ +package link + +import ( + "context" + "encoding/json" + "testing" + "time" + + "github.com/nats-io/nats.go" +) + +// The enrolment round trip against a real server. +// +// **This is the test that keeps a reason from becoming folklore.** The reply address travels in the +// request's payload because a JetStream consumer's delivery has had the transport's reply field +// claimed for its own ack address — which is a fact about a server, not a rule anybody can check by +// reading. Both halves are asserted here: that the field really is eaten, and that the answer +// reaches the node anyway. +// +// docker run -d --rm --name t -p 14223:4222 nats:2.10-alpine -js +// MESH_TEST_NATS=nats://127.0.0.1:14223 go test ./internal/link/ -run TestNatsAnEnrolment + +// asking is a conversation on a bus with no TLS. Built directly rather than through Present because +// the pin is what Present adds and PinnedConfig's own tests cover it; what is under test here is the +// address the answer comes back on. +func asking(t *testing.T, conn *nats.Conn, js nats.JetStreamContext, node string) *natsAsking { + t.Helper() + inbox, err := enrolInbox(node) + if err != nil { + t.Fatal(err) + } + answers, err := conn.SubscribeSync(inbox) + if err != nil { + t.Fatal(err) + } + if err := conn.Flush(); err != nil { + t.Fatal(err) + } + a := &natsAsking{conn: conn, js: js, inbox: inbox, answers: answers, lost: make(chan error, 1)} + t.Cleanup(a.Close) + return a +} + +// theMeshAnswers stands in for the controller: it consumes the enrolment off the stream, reads the +// reply address out of the payload — never from the transport field — and answers there. It reports +// what the transport field actually held, which is the claim design 25 §2 rests on. +func theMeshAnswers(t *testing.T, conn *nats.Conn, js nats.JetStreamContext, + reply EnrolReply) <-chan string { + t.Helper() + sawReplyField := make(chan string, 1) + sub, err := js.Subscribe(EnrolSubject, func(msg *nats.Msg) { + select { + case sawReplyField <- msg.Reply: + default: + } + var addressed struct { + ReplyTo string `json:"reply_to"` + } + if err := json.Unmarshal(msg.Data, &addressed); err != nil || addressed.ReplyTo == "" { + _ = msg.Ack() + return + } + body, _ := json.Marshal(reply) + // Published explicitly to the address the payload named, never msg.Respond — which would + // send it to whatever the transport's reply field holds, and that is the point. + _ = conn.Publish(addressed.ReplyTo, body) + _ = msg.Ack() + }, nats.Durable("controller-standin"), nats.ManualAck()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sub.Unsubscribe() }) + return sawReplyField +} + +// An enrolment is answered on the address the request carried, and the transport's own reply field +// held something else entirely. +func TestNatsAnEnrolmentIsAnsweredOnTheAddressInItsPayload(t *testing.T) { + conn, js := aBus(t) + const node = "joining" + + sawReplyField := theMeshAnswers(t, conn, js, EnrolReply{Accepted: true, Node: node, Password: "p"}) + a := asking(t, conn, js, node) + + request, _ := json.Marshal(EnrolRequest{Node: node, Secret: "t"}) + answer, err := a.Ask(context.Background(), request, 8*time.Second) + if err != nil { + t.Fatalf("no answer reached the node: %v", err) + } + var reply EnrolReply + if err := json.Unmarshal(answer, &reply); err != nil { + t.Fatal(err) + } + if !reply.Accepted || reply.Node != node { + t.Fatalf("the answer was not the mesh's: %+v", reply) + } + + // And the field the answer would have gone to, had it used the transport's: the consumer's own + // ack address. If a future server stopped doing this, the payload-borne address would still + // work and this line is what would say the reason had changed. + select { + case field := <-sawReplyField: + if field == a.inbox { + t.Fatalf("the transport's reply field held this node's inbox (%s), so the payload "+ + "address is no longer load-bearing — check design 25 §2 before relying on it", field) + } + if field == "" { + t.Fatal("the transport's reply field was empty rather than claimed, which is a third " + + "behaviour from the two design 25 §2 describes") + } + case <-time.After(2 * time.Second): + t.Fatal("the stand-in never saw the request") + } +} + +// A request that names no reply address is not answered, and the node says so as a mesh that is not +// running rather than hanging. The controller has nowhere to send an answer, which is the failure +// the payload field exists to make impossible — asserted so that a request built without it fails +// loudly here rather than quietly on a machine. +func TestNatsAnEnrolmentWithNoReplyAddressIsNotAnswered(t *testing.T) { + conn, js := aBus(t) + const node = "silent" + + theMeshAnswers(t, conn, js, EnrolReply{Accepted: true, Node: node}) + a := asking(t, conn, js, node) + + // Published without going through Ask, so the reply address is genuinely absent. + request, _ := json.Marshal(EnrolRequest{Node: node, Secret: "t"}) + if _, err := js.Publish(EnrolSubject, request); err != nil { + t.Fatal(err) + } + if _, err := a.Ask(context.Background(), []byte(`{"node":"`+node+`"}`), 0); err == nil { + t.Fatal("a node with no answer coming was told it had one") + } +} diff --git a/internal/link/enrol.go b/internal/link/enrol.go index 8f755c8..78abc15 100644 --- a/internal/link/enrol.go +++ b/internal/link/enrol.go @@ -6,10 +6,7 @@ import ( "encoding/json" "errors" "fmt" - "net/url" "time" - - amqp "github.com/rabbitmq/amqp091-go" ) // The wire format shared with the control plane, which defines it separately because this binary @@ -53,6 +50,16 @@ type EnrolRequest struct { // on a token this key already spent (novox/hq issue 083). Proof []byte `json:"proof,omitempty"` + // ReplyTo is where the mesh's answer goes, as a field of the request rather than the + // transport's own reply address (design 25 §2). Written by the transport that needs it — + // withReplyTo, once per attempt — because a request going into a stream has had the transport's + // reply field claimed for the consumer's ack address before the controller ever reads it. + // + // Empty on the bus the mesh runs on today, where the delivery carries the reply queue and the + // field means what it has always meant. Named here so both sides of the wire hold the same + // field name, which is what the shape test on each side is for. + ReplyTo string `json:"reply_to,omitempty"` + // Tunnel is the tunnel this node found and whose key it took as its overlay key (novox/hq ADR // 0105): everything about it but that key. Sent with the keys because it is one of them — // OverlayKey above IS this tunnel's public key when this is set — and the mesh composes the @@ -147,52 +154,15 @@ func answered(reply EnrolReply, asking time.Duration) (again bool, err error) { // was issued and the secret is its password. So this is not how the node gets in — it is what it // says once it is in, and the secret travels again because the control plane must not have to ask // the broker who connected. -func Enrol(ctx context.Context, address, pin, node, secret string, public []byte, +func Enrol(ctx context.Context, to Approach, node, secret string, public []byte, overlayKey, sealingKey, servingKey string, profile map[string]any, proof []byte, tunnel *Tunnel, timeout time.Duration) (EnrolReply, error) { - config, err := PinnedConfig(pin) - if err != nil { - return EnrolReply{}, 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), address) - - conn, err := amqp.DialConfig(dsn, amqp.Config{ - TLSClientConfig: config, - Dial: amqp.DefaultDial(timeout), - }) - if err != nil { - if errors.Is(err, ErrWrongCertificate) { - return EnrolReply{}, err - } - // Not quoted back: the DSN carries the one-time secret. - return EnrolReply{}, fmt.Errorf("cannot reach the broker at %s as %s: %w", address, node, err) - } - defer conn.Close() - - channel, err := conn.Channel() - if err != nil { - return EnrolReply{}, err - } - defer channel.Close() - - // 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 { - return EnrolReply{}, fmt.Errorf( - "cannot declare this node's queue %s: %w", QueueFor(node), err) - } - - replies, err := channel.Consume(queue.Name, "", true, false, false, false, nil) + asking, err := Present(ctx, to, node, secret, timeout) if err != nil { return EnrolReply{}, err } + defer asking.Close() request := EnrolRequest{Node: node, Secret: secret, PublicKey: public, OverlayKey: overlayKey, SealingKey: sealingKey, ServingKey: servingKey, Profile: profile, @@ -202,81 +172,47 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte return EnrolReply{}, err } - // Asked, and asked again with the same request while the mesh says "try again": the keys - // this node generated are the ones it keeps, so the same request is the same enrolment, and - // the mesh holds the token for it (novox/hq issue 083). - ask := func() (string, error) { - correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano()) - publish, cancel := context.WithTimeout(ctx, timeout) - defer cancel() - if err := channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false, - amqp.Publishing{ - ContentType: "application/json", - CorrelationId: correlation, - ReplyTo: queue.Name, - Body: body, - }); err != nil { - return "", fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err) - } - return correlation, nil - } - correlation, err := ask() - if err != nil { - return EnrolReply{}, err - } + // Asked, and asked again with the same request while the mesh says "try again": the keys this + // node generated are the ones it keeps, so the same request is the same enrolment, and the mesh + // holds the token for it (novox/hq issue 083). began := time.Now() - - // 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(timeout) - defer deadline.Stop() - closed := conn.NotifyClose(make(chan *amqp.Error, 1)) - for { + answer, err := asking.Ask(ctx, body, timeout) + if err != nil { + return EnrolReply{}, err + } + var reply EnrolReply + if err := json.Unmarshal(answer, &reply); err != nil { + return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err) + } + again, err := answered(reply, time.Since(began)) + if err != nil { + return reply, err + } + if !again { + return reply, nil + } select { case <-ctx.Done(): return EnrolReply{}, ctx.Err() - case reason := <-closed: - return EnrolReply{}, fmt.Errorf("the broker closed the connection: %v", reason) - case <-deadline.C: - return EnrolReply{}, 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", timeout) - case delivery, ok := <-replies: - if !ok { - return EnrolReply{}, errors.New("the broker stopped delivering") - } - // Anything else on this queue is not the answer to this question. - if delivery.CorrelationId != correlation { - continue - } - var reply EnrolReply - if err := json.Unmarshal(delivery.Body, &reply); err != nil { - return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err) - } - again, err := answered(reply, time.Since(began)) - if err != nil { - return reply, err - } - if !again { - return reply, nil - } - select { - case <-ctx.Done(): - return EnrolReply{}, ctx.Err() - case <-time.After(AskAgainAfter): - } - if correlation, err = ask(); err != nil { - return EnrolReply{}, err - } - if !deadline.Stop() { - select { - case <-deadline.C: - default: - } - } - deadline.Reset(timeout) + case <-time.After(AskAgainAfter): } } } + +// withReplyTo writes this attempt's reply address into the request, as a field of its own. +// +// **Written into the bytes rather than carried beside them**, because the whole point is that the +// address survives a stream: a JetStream consumer's delivery has had the transport's reply field +// claimed for its own ack address, so a reply address that is not in the payload is one the +// controller cannot read (design 25 §2). Done by decoding and re-encoding rather than by setting the +// field before marshalling, so one request can be asked again with a fresh address each time without +// the caller knowing that is what happens. +func withReplyTo(request []byte, inbox string) ([]byte, error) { + var fields map[string]any + if err := json.Unmarshal(request, &fields); err != nil { + return nil, fmt.Errorf("this node's own enrolment request cannot be read back: %w", err) + } + fields["reply_to"] = inbox + return json.Marshal(fields) +} diff --git a/internal/link/hearing_nats.go b/internal/link/hearing_nats.go index e0f90e5..d3416be 100644 --- a/internal/link/hearing_nats.go +++ b/internal/link/hearing_nats.go @@ -28,6 +28,11 @@ import ( // answer to novox/hq issue 107. What the drain in run.go still answers is the live case: three // pushes to a *connected* node are three deliveries whatever the stream later retains. +// EnrolSubject is where a joining machine asks. One subject for every node, because a machine +// enrolling has no name the mesh has agreed to yet — which is why its authority to publish here is +// the whole of what its enrolment user may do. +const EnrolSubject = "mesh.control.enrol" + // DeclareSubject is where this node's declaration lands. Its own, and no other node's: a host's // account subscribes exactly this and the subject is the authority on which node a declaration is // for.