The other implementation behind the seam, so the store-window guarantee now has both: one loop, one message at a time, the same window deciding. What differs is where a held message lives, and that is the whole point of the move — the AMQP side keeps an unacknowledged delivery in this process, bounded by the prefetch and lost if the controller stops; this keeps eight bytes saying when the window opened, and the message stays the server's. Checked against a running server, seven claims that reasoning cannot answer: a report is heard and leaves the work queue; one the store cannot take is naked with a delay, stays in the stream, and is recorded when the store returns; one about a superseded declaration is settled without being acted on; one the store never takes is let go once the bound passes; a heartbeat is heard and nothing is persisted; and the enrolment answer reaches the address the request carried in its payload — the test design 25 §2 asks for, so the reason for that field cannot quietly become folklore. Three things the wiring forced into the open: **The controller could not have consumed a module event.** Its permissions granted no event subject to subscribe and no ack subject on the events stream, so every announcement would have been redelivered for ever, refused by the list it already had. Both narrow: each followed subject named, not `mesh.mod.*.>`. **The controller's consumers are not derived.** It files no manifest, so its authority cannot come from a declaration that does not exist; they sit beside the mesh's own streams and are asserted the same way. No max-deliver on CONTROL — the window's bound is the controller's, and a server that dead-lettered first would discard the push the stream exists to protect. **Channels, not callbacks.** The library would run a handler on its own goroutine, and the window's bookkeeping is unlocked because the AMQP loop never had two.
286 lines
10 KiB
Go
286 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.
|
|
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:
|
|
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
|
|
}
|
|
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)
|
|
}
|
|
|
|
// About is nothing here, and that is the point.
|
|
//
|
|
// Setting a held message aside when a newer one about the same thing arrives is what a controller
|
|
// holding deliveries in memory can do. A naked message belongs to the server and comes back
|
|
// whatever happened meanwhile, so the question "is this the past?" is answered by what the message
|
|
// says instead — the digest of the declaration a report is about (window.go, design 25 §3).
|
|
func (m *natsControl) About(string) {}
|
|
|
|
// 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"`
|
|
}
|