Files
mesh-host/internal/link/hearing_nats.go
T
jochen 5e189c2fca
mesh/merge-gate pass: builds mesh-host → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 of 4 compose)
mesh/repo-check pass: THE CHANGE ALTERS ITS OWN CHECK (merge-check.sh): main's version judged it; the change's judges the pull requests after it merges; it…
mesh/delivery delivered
Say a tool check's errors once and as they stand, and test the link on a real bus
The evidence read "asking it: asking <subject>": the words are now the link's own.
A deadline is said as the time the check gave it. A refusal is said only when the bus
refused this question, not an earlier one under a grant since widened. merge-check.sh
runs the tests that need a bus against a throwaway nats-server when the toolchain has
one, and says so when it does not (hq issue 331).
2026-10-08 17:01:15 +02:00

263 lines
9.7 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" }
// AskReportSubject is where the mesh asks this node to say again what it last applied (novox/hq to-be 45
// §6): its own, on core NATS and off any stream. The answer is an ordinary report on its report subject,
// the one thing it already says — nobody's inbox is answered.
func AskReportSubject(node string) string { return "mesh.node." + node + ".ask.report" }
// 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
asking *nats.Subscription
asked chan struct{}
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),
// 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.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),
asked: make(chan struct{}, 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
l.hearAsks()
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
}
// hearAsks listens for the mesh asking what this node last applied (the `report` verb). One question
// waiting is enough: a second asked before the first is taken is the same question. A bus whose user
// list is older than the grant refuses the subscription asynchronously — said by the client, and the
// mesh then sends again instead of asking; nothing here depends on being asked.
func (l *natsLink) hearAsks() {
if l.asked == nil {
l.asked = make(chan struct{}, 1)
}
asking, err := l.conn.Subscribe(AskReportSubject(l.node), func(*nats.Msg) {
select {
case l.asked <- struct{}{}:
default:
}
})
if err == nil {
l.asking = asking
}
}
func (l *natsLink) Declarations() <-chan Declaration { return l.arrived }
func (l *natsLink) AskedToReport() <-chan struct{} { return l.asked }
func (l *natsLink) Lost() <-chan error { return l.lost }
func (l *natsLink) Close() {
if l.sub != nil {
_ = l.sub.Unsubscribe()
}
if l.asking != nil {
_ = l.asking.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)
}
func (l *natsLink) Health(ctx context.Context, node string, body []byte) error {
return OverNATS{Conn: l.conn, JS: l.js}.Health(ctx, node, body)
}
// Ask asks a module's tool for a declared tool check (novox/hq ADR 0240, to-be 48 §3). The link the
// node-engine holds is this one, not OverNATS: without it here the queue's ToolBus assertion failed on
// every live link and each tool check read "no link to the bus is open" while the link was up.
func (l *natsLink) Ask(ctx context.Context, subject string, body []byte) ([]byte, error) {
before := l.conn.LastError()
reply, err := OverNATS{Conn: l.conn, JS: l.js}.Ask(ctx, subject, body)
switch {
case err == nil:
return reply, nil
case errors.Is(err, nats.ErrNoResponders):
return nil, fmt.Errorf("nothing on this machine answers %s", subject)
}
if refused := l.refusal(subject, before); refused != "" {
return nil, fmt.Errorf("%s: %v%s", subject, err, refused)
}
if errors.Is(err, context.DeadlineExceeded) || errors.Is(err, nats.ErrTimeout) {
// Said as the deadline, which the caller words with the time it gave the question.
return nil, fmt.Errorf("no answer from %s: %w", subject, context.DeadlineExceeded)
}
return nil, fmt.Errorf("%s: %w", subject, err)
}
// What the queue asks of a live link beside Bus, each by a type assertion that fails quietly: said here
// so a link that drops one does not build.
var (
_ Bus = (*natsLink)(nil)
_ HealthBus = (*natsLink)(nil)
_ ToolBus = (*natsLink)(nil)
_ Asker = (*natsLink)(nil)
_ Asked = (*natsLink)(nil)
)
// 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() }