Files
mesh-host/internal/link/hearing_nats.go
T
jochen 1fc1e74cc3
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: its merge-check.sh passed
mesh/delivery delivered
mesh/delivery-group group feat/a-module-says-how-it-is-healthy delivered: every member is delivered
Judge whether what a module runs stays up, and say it (hq ADR 0240, to-be 48 Phase A)
A container that crash-looped after its compose applied passed every check the
gate had: nothing looked at what a module runs. The node-engine now judges every
long-running resource on every look — one read of the runtime, one per service
manager — keeps the restarts it counts across recreates and its own restarts,
and says the state in every report and as an event on change, again every minute
while not healthy. It reads only; nothing is restarted for being unhealthy.
2026-10-07 02:28:15 +02:00

231 lines
8.3 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)
}
// 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() }