Files
mesh-controller/internal/link/receive_nats.go
T
jschoubben aecac5bda2 One bus: the AMQP transport is gone from the controller
The mesh runs on the seat's bus alone (novox/hq ADR 0131, design 28 task 5.5). The old
transport's consume loop, build request, tool ask, management API and account scoping are
deleted, and the bus switch with them; the controller connects to the broker seat and to
nothing else. The store-window tests keep their assertions on a bus-less fake, and the tests
that only made sense for the old transport's in-memory holding go with it.
2026-09-28 03:36:16 +02:00

289 lines
10 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
)
// The consume side on the bus being built.
//
// The shape is the AMQP one's, because the seam made them comparable: one loop, one message at a
// time, and the same window deciding. What differs is where a held message lives — and that is the
// whole point of the move. On the bus the mesh has, holding one means keeping an unacknowledged
// delivery in this process, bounded by the prefetch and lost if the controller stops. Here it is a
// `nak` with a delay: the message stays the server's, the controller keeps nothing but the moment
// it first could not take it, and a controller that restarts mid-window has nothing to lose.
// natsInbound consumes what nodes and modules say over NATS.
// Prefetch is how many controls the controller holds unacknowledged at once: the store window
// (ADR 0083) is the server's, and this is what it may hand this process ahead of its acting.
const Prefetch = 64
type natsInbound struct {
js *broker.JetStream
// follows is the kinds asked for beyond what nodes say (Also). The events those are are the
// only ones the controller subscribes, and only when something is listening.
follows map[string]bool
// since is when the controller first could not take a message, by that message's place in its
// stream.
//
// **A timestamp, not a message.** This is the whole difference the move buys: the AMQP side
// keeps the delivery, and this keeps eight bytes saying when the window opened. A controller
// that restarts loses these and starts the window again, which is correct — it is holding
// nothing, and the messages are all still on the server.
since map[uint64]time.Time
}
// Nats is the consume side of the bus being built.
func Nats(js *broker.JetStream) Inbound {
return &natsInbound{js: js, follows: map[string]bool{}, since: map[uint64]time.Time{}}
}
// Also records one more kind to subscribe. Nothing is subscribed here: the controller's consumer on
// the events stream carries both of these as filters, so it is created once, in Receive, with
// whatever was asked for — and not at all when nothing was.
func (n *natsInbound) Also(kind string) error {
switch kind {
case KindModuleMoved, KindCatchUp, KindSourceMoved:
n.follows[kind] = true
return nil
default:
return fmt.Errorf("nothing subscribes %s separately on this bus", kind)
}
}
func (n *natsInbound) Close() {}
// Receive consumes until the context ends.
//
// Three subscriptions, and each is a channel the one loop selects on. **Channels rather than
// callbacks**: the library would run a handler on its own goroutine, and the window's bookkeeping —
// which message is held, and since when — is read and written without a lock because the AMQP loop
// never had two. A second goroutine would make that wrong in a way no test would catch.
func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error {
if err := broker.AssertMeshConsumers(n.js); err != nil {
return err
}
js, conn := n.js.Context(), n.js.Conn()
// What nodes say, off the CONTROL stream. Bound to the durable the controller asserted rather
// than creating one here: the consumer is an object with a configuration — ack policy, ack
// wait, redelivery — and a client that creates its own would be a second opinion about it.
control := make(chan *nats.Msg, Prefetch)
said, err := js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName))
if err != nil {
return fmt.Errorf("subscribing to what nodes say: %w", err)
}
defer func() { _ = said.Unsubscribe() }()
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
// they are their own guarantee: a lost one is the next one.
beats := make(chan *nats.Msg, Prefetch)
alive, err := conn.ChanSubscribe(AliveSubjects, beats)
if err != nil {
return fmt.Errorf("subscribing to heartbeats: %w", err)
}
defer func() { _ = alive.Unsubscribe() }()
// The events the controller follows, when something is listening for them.
var events chan *nats.Msg
if len(n.follows) > 0 {
events = make(chan *nats.Msg, Prefetch)
followed, err := js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName))
if err != nil {
return fmt.Errorf("subscribing to what the catalogue says: %w", err)
}
defer func() { _ = followed.Unsubscribe() }()
}
// A connection that dropped is said, not discovered. A controller whose bus connection is gone
// is a mesh where nothing can be told anything.
gone := make(chan error, 1)
conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) {
select {
case gone <- err:
default:
}
})
for {
select {
case <-ctx.Done():
return nil
case err := <-gone:
return fmt.Errorf("the bus connection dropped: %w", err)
case msg := <-beats:
n.deliver(ctx, act, msg, false)
case msg := <-events:
n.deliver(ctx, act, msg, true)
case msg, ok := <-control:
if !ok {
return errors.New("the bus stopped delivering")
}
n.deliver(ctx, act, msg, true)
}
}
}
// deliver names one message and hands it to the loop, or drops it where the mesh has no name for
// its subject — which cannot happen through a filter the controller wrote, and is said rather than
// ignored for exactly that reason.
func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control),
msg *nats.Msg, streamed bool) {
kind, known := kindOfSubject(msg.Subject)
if !known {
if streamed {
_ = msg.Term()
}
return
}
m := &natsControl{kind: kind, msg: msg, on: n}
if streamed {
// A message with no metadata is not from a stream, whatever it was delivered on, and the
// window has nothing to hold it by. Said by leaving the sequence at zero.
if meta, err := msg.Metadata(); err == nil {
m.seq = meta.Sequence.Stream
m.delivered = meta.NumDelivered
}
}
act(ctx, m)
}
// kindOfSubject is how this transport's addressing becomes what the mesh calls a message.
//
// By subject, which is the only thing the server enforces: a body claiming to be a report does not
// make it one, and on this bus the subject an account may publish *is* its authority (design 29
// §2). The mirror of the routing-key table on the bus the mesh has.
func kindOfSubject(subject string) (string, bool) {
switch subject {
case EnrolSubject:
return KindEnrolment, true
case BuiltSubject:
return KindBuilt, true
}
if node, rest, ok := strings.Cut(strings.TrimPrefix(subject, "mesh.control."), "."); ok &&
node != "" && !strings.Contains(node, ".") {
switch rest {
case "report":
return KindReport, true
case "alive":
return KindHeartbeat, true
}
}
switch subject {
case broker.ControllerFollows[0]:
return KindModuleMoved, true
case broker.ControllerFollows[1]:
return KindCatchUp, true
case broker.ControllerFollows[3]:
return KindSourceMoved, true
case BuildOutcome():
// A build's outcome is the role's event now, so it arrives on the events stream rather than
// the control branch — and is acted on by the same handler, because what the controller does
// with it did not change (novox/hq ADR 0121).
return KindBuilt, true
}
return "", false
}
// natsControl is one message from the bus being built, as the controller reads it.
type natsControl struct {
kind string
msg *nats.Msg
on *natsInbound
// seq is this message's place in its stream; zero for a core message, which has none and
// cannot be held.
seq uint64
// delivered is how many times the server has handed this message over, this time included.
delivered uint64
}
func (m *natsControl) Kind() string { return m.kind }
func (m *natsControl) Body() []byte { return m.msg.Data }
// Redelivered is what the server counted, not what the controller remembers. Which is the answer to
// a question the AMQP side could only guess at across a restart: an enrolment redelivered because
// the controller stopped mid-answer reads as redelivered to the controller that comes back.
func (m *natsControl) Redelivered() bool { return m.delivered > 1 }
func (m *natsControl) HeldFor() time.Duration {
if m.seq == 0 {
return 0
}
first, held := m.on.since[m.seq]
if !held {
return 0
}
return time.Since(first)
}
// Answer publishes to the reply subject the request carries **in its payload**.
//
// Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had
// that field claimed for the consumer's own ack address, so answering it would send the reply to
// `$JS.ACK.CONTROL.controller.…` and the enrolling node would wait out its timeout. Verified
// against a running server (design 25 §2), which is why it is a field of the request and this reads
// it from there.
func (m *natsControl) Answer(ctx context.Context, body []byte) error {
var addressed replyAddressed
if err := json.Unmarshal(m.msg.Data, &addressed); err != nil {
return fmt.Errorf("that request cannot be read, so its reply address cannot be: %w", err)
}
if addressed.ReplyTo == "" {
return errors.New("that request named no reply subject in its payload, so nothing can be " +
"told the answer")
}
return m.on.js.Conn().PublishMsg(&nats.Msg{Subject: addressed.ReplyTo, Data: body})
}
func (m *natsControl) Took() error {
m.forget()
if m.seq == 0 {
// Core NATS: nothing is keeping it, so there is nothing to settle.
return nil
}
return m.msg.Ack(nats.Context(context.Background()))
}
// Drop terminates the delivery: understood, and the server is told not to send it again. Different
// from an ack only in the server's own accounting, which is where somebody asking "what happened to
// that message" will look.
func (m *natsControl) Drop() error {
m.forget()
if m.seq == 0 {
return nil
}
return m.msg.Term()
}
// Hold hands the message back with a delay, and remembers when the window opened.
func (m *natsControl) Hold(after time.Duration) error {
if m.seq == 0 {
return errors.New("a message that is not in a stream cannot be held: nothing is keeping it")
}
if _, already := m.on.since[m.seq]; !already {
m.on.since[m.seq] = time.Now()
}
return m.msg.NakWithDelay(after)
}
func (m *natsControl) forget() {
if m.on != nil && m.seq != 0 {
delete(m.on.since, m.seq)
}
}
// replyAddressed is the one field every message that expects an answer carries.
type replyAddressed struct {
ReplyTo string `json:"reply_to,omitempty"`
}