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.
188 lines
6.5 KiB
Go
188 lines
6.5 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// The host's link on the bus being built.
|
|
//
|
|
// Two things a host does here that it cannot do on the other bus, and one it must not try.
|
|
//
|
|
// **It declares nothing.** On the bus the mesh has, a host declares its own queue on connecting,
|
|
// because a queue that is not there means a node that hears nothing. Here the object it reads
|
|
// through is a durable consumer, and a host's account reaches no part of the JetStream API — by
|
|
// design, because the controller is the only writer of consumer definitions (design 25 §3). So the
|
|
// host **binds** to a consumer the controller made when this node enrolled, and a missing one is
|
|
// said as what it is rather than quietly created with whatever configuration this client happens to
|
|
// default to.
|
|
//
|
|
// **It gets order for free, and keeps the drain anyway.** The declaration subject is last-per-subject
|
|
// (design 29 §4), so a node that was away receives exactly the current declaration rather than a
|
|
// queue of superseded ones, and the stream's sequence orders them definitively — the wire-level
|
|
// 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.
|
|
func DeclareSubject(node string) string { return "mesh.node." + node + ".declare" }
|
|
|
|
// natsURL is a bus address as the client wants it. A membership records host and port, because that
|
|
// is what genesis sealed into it and what the other transport takes; the scheme is this transport's
|
|
// own business.
|
|
func natsURL(address string) string {
|
|
if strings.Contains(address, "://") {
|
|
return address
|
|
}
|
|
return "nats://" + address
|
|
}
|
|
|
|
// natsLink is this node's connection as a JetStream subscription.
|
|
type natsLink struct {
|
|
conn *nats.Conn
|
|
js nats.JetStreamContext
|
|
sub *nats.Subscription
|
|
node string
|
|
arrived chan Declaration
|
|
lost chan error
|
|
}
|
|
|
|
func dialNats(ctx context.Context, m Membership, timeout time.Duration) (Link, error) {
|
|
// Pinned exactly as the other transport is, and for once the Go client makes that easy: it
|
|
// takes a *tls.Config, so the same PinnedConfig with the same VerifyPeerCertificate does the
|
|
// work. **The constraint recorded against the tool runtime does not apply here** — that client
|
|
// takes PEM strings with no verify hook, which is why the bus's certificate must carry a name
|
|
// matching the address *modules* dial it by. A host checks the fingerprint and nothing else.
|
|
config, err := PinnedConfig(m.Fingerprint)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
opts := []nats.Option{
|
|
nats.Secure(config),
|
|
nats.UserInfo(m.Node, m.Password),
|
|
nats.Name("mesh-host/" + m.Node),
|
|
nats.Timeout(timeout),
|
|
// A node that has silently lost its route notices, rather than holding a connection the
|
|
// server forgot about and believing it is still in the mesh.
|
|
nats.PingInterval(10 * time.Second),
|
|
nats.MaxPingsOutstanding(2),
|
|
// Reconnection is the caller's: Hold already decides when to try again and how long to
|
|
// wait, and a client quietly reconnecting underneath it would make that reasoning a
|
|
// duplicate of the library's.
|
|
nats.NoReconnect(),
|
|
}
|
|
|
|
conn, err := nats.Connect(natsURL(m.Broker), opts...)
|
|
if err != nil {
|
|
if errors.Is(err, ErrWrongCertificate) {
|
|
return nil, err
|
|
}
|
|
return nil, fmt.Errorf("cannot reach the bus at %s: %w", m.Broker, err)
|
|
}
|
|
js, err := conn.JetStream()
|
|
if err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("the bus at %s has no JetStream: %w", m.Broker, err)
|
|
}
|
|
|
|
l := &natsLink{
|
|
conn: conn, js: js, node: m.Node,
|
|
arrived: make(chan Declaration, drainDepth),
|
|
lost: make(chan error, 1),
|
|
}
|
|
|
|
// Bound to the consumer the controller made for this node, named after the node because that is
|
|
// what the node's own ack grant allows (`$JS.ACK.NODES.<node>.>`).
|
|
feed := make(chan *nats.Msg, drainDepth)
|
|
// The subject as well as the binding: the client checks what is asked for against the
|
|
// consumer's own filter, and an empty subject is refused rather than taken to mean "whatever
|
|
// that consumer delivers".
|
|
sub, err := js.ChanSubscribe(DeclareSubject(m.Node), feed, nats.Bind("NODES", m.Node))
|
|
if err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf(
|
|
"this node cannot read its declarations: %w. The mesh creates that when a node enrols, "+
|
|
"and a host may not create one itself — so this is the mesh's to answer, not this "+
|
|
"machine's", err)
|
|
}
|
|
l.sub = sub
|
|
|
|
conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) {
|
|
select {
|
|
case l.lost <- fmt.Errorf("the link dropped: %w", err):
|
|
default:
|
|
}
|
|
})
|
|
conn.SetClosedHandler(func(*nats.Conn) {
|
|
select {
|
|
case l.lost <- errors.New("the link closed"):
|
|
default:
|
|
}
|
|
})
|
|
|
|
go func() {
|
|
defer close(l.arrived)
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case msg, ok := <-feed:
|
|
if !ok {
|
|
select {
|
|
case l.lost <- errors.New("the bus stopped delivering"):
|
|
default:
|
|
}
|
|
return
|
|
}
|
|
select {
|
|
case l.arrived <- natsDeclaration{msg}:
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}()
|
|
|
|
return l, nil
|
|
}
|
|
|
|
func (l *natsLink) Declarations() <-chan Declaration { return l.arrived }
|
|
func (l *natsLink) Lost() <-chan error { return l.lost }
|
|
|
|
func (l *natsLink) Close() {
|
|
if l.sub != nil {
|
|
_ = l.sub.Unsubscribe()
|
|
}
|
|
if l.conn != nil {
|
|
l.conn.Close()
|
|
}
|
|
}
|
|
|
|
func (l *natsLink) Report(ctx context.Context, node string, body []byte) error {
|
|
return OverNATS{Conn: l.conn, JS: l.js}.Report(ctx, node, body)
|
|
}
|
|
|
|
func (l *natsLink) Alive(ctx context.Context, node string, body []byte) error {
|
|
return OverNATS{Conn: l.conn, JS: l.js}.Alive(ctx, node, body)
|
|
}
|
|
|
|
// natsDeclaration is one declaration off the NODES stream.
|
|
type natsDeclaration struct{ msg *nats.Msg }
|
|
|
|
func (d natsDeclaration) Body() []byte { return d.msg.Data }
|
|
|
|
// Handled acknowledges it. The ack goes to this node's own ack subject, which is the one thing
|
|
// besides its reports a node's account may publish.
|
|
func (d natsDeclaration) Handled() error { return d.msg.Ack() }
|