novox/hq 04-ISSUES/146, the layers behind the three already fixed. A new membership says which bus it is for. Empty meant 'whatever the mesh runs today' while two buses existed, and became a refusal the moment one did: an enrolled node came up and reconnected for ever against its own record. The enrolling client takes its inboxes in the space its user may listen in. A JetStream publish waits for the stream's acknowledgement on an inbox the client picks, and its default is one this user may not subscribe to — so the enrolment failed with a permissions violation on a subject nobody had chosen. And the enrolment publish carries a message id, so the client's own retry is discarded by the stream rather than enrolling the machine twice. That one is not finished: the duplicate survives it, and the issue says where the trail stops.
187 lines
7.1 KiB
Go
187 lines
7.1 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"crypto/sha256"
|
|
"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.
|
|
// **Its own inbox space, because that is the only one it may listen in** (novox/hq
|
|
// 04-ISSUES/146). A JetStream publish waits for the stream's acknowledgement on an inbox the
|
|
// client picks, and the client's default is `_INBOX.<random>` — which this user may not
|
|
// subscribe to, so the enrolment failed with a permissions violation on a subject nobody had
|
|
// chosen. The permission is `_INBOX.enrol.<node>.>` (design 25 §6), so the client is told to
|
|
// pick its inboxes there; the reply address below is in the same space for the same reason.
|
|
conn, err := nats.Connect(natsURL(to.Address),
|
|
nats.Secure(config),
|
|
nats.CustomInboxPrefix("_INBOX.enrol."+node),
|
|
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.
|
|
//
|
|
// **Once, however many times it is sent** (novox/hq 04-ISSUES/146). The client re-publishes when
|
|
// an acknowledgement is slow, and the mesh enrolled the machine on each copy — minting a second
|
|
// credential, which replaced the first, which is the one the node had already been given. The
|
|
// machine then reconnected for ever as a user whose password the mesh had rotated out from under
|
|
// it, and the controller's log said "enrolled anchor" twice in the same second.
|
|
//
|
|
// The id is the message: the same bytes carry the same id, so the stream discards the client's
|
|
// own retry, and a genuine second attempt — which carries a new reply address — is a different
|
|
// message and is let through.
|
|
sum := sha256.Sum256(addressed)
|
|
if _, err := a.js.Publish(EnrolSubject, addressed,
|
|
nats.MsgId(hex.EncodeToString(sum[:])), 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)
|
|
}
|
|
}
|
|
}
|