The last of the host's link that still named a transport. `Asking` is one enrolment conversation — a connection made with the token, a question asked, and an answer waited for — and it is its own seam rather than part of `Link` because almost nothing about it is the same: the credential is a one-time secret, there is no declaration to hear, and 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. `Enrol`'s thirteen arguments became an `Approach` — where, which certificate, which bus — and the request it already had. The token says nothing about which bus, and does not need to: every token names the one the mesh runs on today until the rollout. **The reply address is the whole of what changes on the new bus**, and it is forced rather than preferred. Verified against a running server, both halves: the answer reaches the node at the address its request carried in the payload, and the transport's own reply field held something else entirely by the time the consumer saw it — the consumer's ack address, exactly as design 25 §2 says. The test asserts the field is *not* the node's inbox, so a future server that stopped claiming it would fail this rather than let the reason quietly become folklore. The inbox is under `_INBOX.enrol.<node>.`, which is exactly what the enrolling user may subscribe and no wider, with a random tail per attempt: 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. Subscribed before anything is published, because a node that published first could miss an answer to a question nobody was listening for.
167 lines
5.7 KiB
Go
167 lines
5.7 KiB
Go
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.<node>.`, 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.<node>`, 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)
|
|
}
|
|
}
|
|
}
|