The consume side on NATS, and the window held by the server

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.
This commit is contained in:
2026-09-27 00:53:33 +02:00
parent 06cf3c04e5
commit 88bef39952
9 changed files with 796 additions and 3 deletions
+7
View File
@@ -28,6 +28,13 @@ type Consumer struct {
// Queue is the queue group, set for a seat's worker so that "exactly one holder" survives a // Queue is the queue group, set for a seat's worker so that "exactly one holder" survives a
// seat later being relaxed to several. Authority and delivery are kept separate on purpose. // seat later being relaxed to several. Authority and delivery are kept separate on purpose.
Queue string Queue string
// Push asks the server to deliver to a subject rather than wait to be pulled.
//
// For the mesh's own consumer, where the controller wants every message to arrive in the one
// loop it already runs: pulling would mean a second goroutine fetching batches and handing
// them over, and a loop that acts on one message at a time is the property the store window
// depends on. A queue group implies this, because a group has nothing to pull from.
Push bool
// AckWaitSeconds before an unacknowledged delivery is redelivered. // AckWaitSeconds before an unacknowledged delivery is redelivered.
AckWaitSeconds int AckWaitSeconds int
// MaxDeliver before the message is dead-lettered; zero for the mesh's default. // MaxDeliver before the message is dead-lettered; zero for the mesh's default.
+9 -1
View File
@@ -39,6 +39,14 @@ func Dial(url string, opts ...nats.Option) (*JetStream, error) {
return &JetStream{conn: conn, js: js}, nil return &JetStream{conn: conn, js: js}, nil
} }
// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat,
// a tool call — where a lost message is answered by the next one or by a timeout the caller
// already handles (design 25 §3).
func (j *JetStream) Conn() *nats.Conn { return j.conn }
// Context is the JetStream handle, for subscribing to what the consumers above define.
func (j *JetStream) Context() nats.JetStreamContext { return j.js }
func (j *JetStream) Close() { func (j *JetStream) Close() {
if j.conn != nil { if j.conn != nil {
j.conn.Close() j.conn.Close()
@@ -112,7 +120,7 @@ func (j *JetStream) EnsureConsumer(c Consumer) error {
// A queue group needs a delivery subject: a pull consumer has no group, and declaring one // A queue group needs a delivery subject: a pull consumer has no group, and declaring one
// without the other is refused by the server with a message that does not say which half is // without the other is refused by the server with a message that does not say which half is
// missing. // missing.
if c.Queue != "" { if c.Queue != "" || c.Push {
want.DeliverSubject = "_DELIVER." + c.Name want.DeliverSubject = "_DELIVER." + c.Name
} }
+10
View File
@@ -143,6 +143,16 @@ func PermissionsFor(p Principal) (Permissions, error) {
pub = []string{"mesh.control.>", "mesh.node.>", "mesh.build.>", "$JS.API.>"} pub = []string{"mesh.control.>", "mesh.node.>", "mesh.build.>", "$JS.API.>"}
sub = []string{"mesh.control.>", "mesh.build.>", "$JS.API.>"} sub = []string{"mesh.control.>", "mesh.build.>", "$JS.API.>"}
// The two events it reacts to, and its ack subject on the stream they arrive from
// (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the
// controller a subscriber to every event in the mesh, and its permission list would stop
// saying what it is for. The ack grant below is scoped per stream because the controller's
// consumer name is the same on both and `$JS.ACK.CONTROL.controller.>` does not cover a
// delivery from EVENTS — a consumer that cannot ack has every message redelivered for
// ever, refused by the list it already has.
sub = append(sub, ControllerFollows...)
pub = append(pub, "$JS.ACK.EVENTS."+ControllerName+".>")
case KindPerson: case KindPerson:
// Tools, and nothing else. Every subject a person may publish is a tool call; a person // Tools, and nothing else. Every subject a person may publish is a tool call; a person
// who could publish an event would be able to claim a module said something. // who could publish an event would be able to claim a module said something.
+76
View File
@@ -150,3 +150,79 @@ func Overlaps() []string {
sort.Strings(clashes) sort.Strings(clashes)
return clashes return clashes
} }
// The mesh's own consumers.
//
// A seat's streams and a module's consumers are derived from declarations (derived.go). These two
// are not: **the controller is not a module and files no manifest**, so its authority and its
// subscriptions cannot come from a declaration that does not exist. They are named here, where the
// mesh's own streams are named, and narrowly — a controller subscribing `mesh.mod.*.event.>` would
// hear every event in the mesh, which it has no business doing and which would make its permission
// list stop explaining anything.
// ControllerName is the controller's durable consumer on each stream it reads, and the name its
// ack subject is derived from (nats.go: `$JS.ACK.<stream>.controller.>`).
const ControllerName = "controller"
// ControllerFollows are the events the controller reacts to: the catalogue saying a module's
// current version moved, and a catalogue that has just started saying it may have missed builds.
//
// **These carry the local names the manifests hold today**, which still spell an event the way a
// routing key on the bus the mesh has does — `module.<module>.<verb>` rather than design 29's bare
// verb — so the derived subject names the module twice. It is consistent, and it is what the
// catalogue actually publishes, so it is what the controller must listen to. It changes when those
// names are converted, and not before: a subscription written against the name design 29 specifies
// would be a controller listening to a subject nothing publishes.
var ControllerFollows = []string{
"mesh.mod.mesh-catalog.event.module.mesh-catalog.upgraded",
"mesh.mod.mesh-catalog.event.module.mesh-catalog.catching-up",
}
// MeshConsumers is what the controller consumes, in the order a person reads it.
//
// **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's,
// not the server's (window.go): a message is held with a nak-and-delay until the controller either
// takes it or gives up and says so. A max-deliver here would dead-letter a push that was being
// held through a store restart — the exact message the stream exists to protect — some minutes
// before the controller had finished deciding about it.
func MeshConsumers() []Consumer {
return []Consumer{
{
Name: ControllerName,
Stream: "CONTROL",
Push: true,
AckWaitSeconds: 30,
Why: "the controller is the single consumer of what nodes say; explicit ack and no " +
"max-deliver, because the store window's bound is the controller's own",
},
{
Name: ControllerName,
Stream: "EVENTS",
Filters: ControllerFollows,
Push: true,
AckWaitSeconds: 30,
MaxDeliver: 5,
Why: "the two events the mesh's own controller reacts to; after max-deliver it " +
"dead-letters, because an announcement it cannot act on will not become actionable",
},
}
}
// Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason
// Asserter is.
type Ensurer interface {
EnsureConsumer(c Consumer) error
}
// AssertMeshConsumers brings the controller's own consumers into being, and says which one failed.
//
// After the streams, necessarily: a consumer on a stream that does not exist is refused, and the
// refusal names the stream rather than the order.
func AssertMeshConsumers(e Ensurer) error {
for _, c := range MeshConsumers() {
if err := e.EnsureConsumer(c); err != nil {
return fmt.Errorf("asserting consumer %s on %s: %w", c.Name, c.Stream, err)
}
}
return nil
}
+2 -2
View File
@@ -20,8 +20,8 @@ accounts {
MESH { MESH {
users = [ users = [
{ user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: {
publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.API.>", "mesh.build.>", "mesh.control.>", "mesh.node.>"] } publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "mesh.build.>", "mesh.control.>", "mesh.node.>"] }
subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.build.>", "mesh.control.>"] } subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.build.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.module.mesh-catalog.catching-up", "mesh.mod.mesh-catalog.event.module.mesh-catalog.upgraded"] }
allow_responses: { max: 1, ttl: "1m" } allow_responses: { max: 1, ttl: "1m" }
} } } }
{ user: "enrolment", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { { user: "enrolment", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: {
+28
View File
@@ -92,6 +92,34 @@ type OverNATS struct {
JS nats.JetStreamContext JS nats.JetStreamContext
} }
// The subjects a node publishes on, and the controller listens to.
//
// One tree, and each name says who it is about: `mesh.control.<node>.…` is a node's own, which is
// what lets a node's account be granted exactly its own prefix and nothing of any other node's
// (design 25 §2, §4). The two that belong to no node — an enrolment, because a machine enrolling
// has no name the mesh has agreed to yet, and a build's outcome, because a builder is not
// reporting about itself — are named directly.
const (
// EnrolSubject is where a joining machine asks. Its enrolment user may publish here and
// nowhere else, so a leaked token buys nothing but the chance to enrol.
EnrolSubject = "mesh.control.enrol"
// BuiltSubject is where a build's outcome lands, for results nobody was waiting for.
BuiltSubject = "mesh.control.built"
// AliveSubjects is every node's heartbeat. Core NATS, never a stream: a lost heartbeat is the
// next heartbeat, and a stream of them is the mesh's least valuable message competing for
// retention with its most valuable (design 25 §3).
AliveSubjects = "mesh.control.*.alive"
)
// ReportSubject is where one node says what it did. On the CONTROL stream, because it is the
// message the store-window guarantee is about (ADR 0083).
func ReportSubject(node string) string { return "mesh.control." + node + ".report" }
// AliveSubject is one node's heartbeat.
func AliveSubject(node string) string { return "mesh.control." + node + ".alive" }
// EventSubject is where a module's event lands. Derived from the emitter, never taken from the // EventSubject is where a module's event lands. Derived from the emitter, never taken from the
// caller: a source that could differ from the subject is an envelope that can lie about its // caller: a source that could differ from the subject is an envelope that can lie about its
// origin, and on NATS the account's permissions make the subject the authority (design 29 §2). // origin, and on NATS the account's permissions make the subject the authority (design 29 §2).
+14
View File
@@ -80,6 +80,20 @@ type EnrolRequest struct {
// it. Nil from a node that found none, which is every converged one. // it. Nil from a node that found none, which is every converged one.
Tunnel *Tunnel `json:"tunnel,omitempty"` Tunnel *Tunnel `json:"tunnel,omitempty"`
// ReplyTo is where the answer goes, as a field of the request rather than the transport's own
// reply address.
//
// **Because a stream eats the transport's field** (design 25 §2, verified against a running
// server): a message a JetStream consumer delivers has had its reply field claimed for that
// consumer's own ack address, so by the time the controller sees an enrolment, the field names
// where the *controller* must acknowledge, not where the node is waiting. An enrolment is the
// case that matters — a caller waiting on an ephemeral inbox, over a subject the store window
// may legitimately delay by several nak cycles.
//
// Empty on the bus the mesh runs on today, where the delivery carries the reply queue and the
// field means what it has always meant.
ReplyTo string `json:"reply_to,omitempty"`
// Redelivered is set by the control plane, never sent: the broker handed this request over a // Redelivered is set by the control plane, never sent: the broker handed this request over a
// second time. Such a request does not finish an enrolment already spent — the first time may // second time. Such a request does not finish an enrolment already spent — the first time may
// have answered, and the node holds what it was told. // have answered, and the node holds what it was told.
+285
View File
@@ -0,0 +1,285 @@
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"`
}
+365
View File
@@ -0,0 +1,365 @@
package link
import (
"context"
"encoding/json"
"errors"
"os"
"sync"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
)
// The consume side against a real server, because what is being checked is what the server does.
//
// Reasoning cannot answer any of these: whether a nak-with-delay really comes back, whether the
// delay is honoured, whether terminating a delivery really stops it, or whether an answer published
// to an address carried in the payload reaches a caller waiting on its own inbox. Each is a claim
// about a server, so each is asked of one:
//
// docker run -d --rm --name t -p 14222:4222 nats:2.10-alpine -js
// MESH_TEST_NATS=nats://127.0.0.1:14222 go test ./internal/link/ -run TestNats
func aBus(t *testing.T) *broker.JetStream {
t.Helper()
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
js, err := broker.Dial(url)
if err != nil {
t.Fatal(err)
}
t.Cleanup(js.Close)
// The mesh's own streams and consumers, asserted the way the controller asserts them — and
// torn down after, so one test's held message is never another's surprise.
for _, s := range broker.MeshStreams() {
_ = js.Context().DeleteStream(s.Name)
}
if err := broker.AssertMeshStreams(js); err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
for _, s := range broker.MeshStreams() {
_ = js.Context().DeleteStream(s.Name)
}
})
return js
}
// serving1 is a controller reading from a real bus, and a way to stop it.
func servingOn(t *testing.T, js *broker.JetStream, l Listener) (*Server, func()) {
t.Helper()
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()},
listener: l, log: quiet()}
ctx, stop := context.WithCancel(context.Background())
done := make(chan struct{})
go func() { defer close(done); _ = s.Serve(ctx) }()
return s, func() {
stop()
<-done
}
}
// counted records reports and can be told to refuse them, from another goroutine.
type counted struct {
mu sync.Mutex
err error
heard []Report
}
func (c *counted) Heard(_ context.Context, r Report) error {
c.mu.Lock()
defer c.mu.Unlock()
if c.err != nil {
return c.err
}
c.heard = append(c.heard, r)
return nil
}
func (c *counted) refusing(err error) {
c.mu.Lock()
defer c.mu.Unlock()
c.err = err
}
func (c *counted) count() int {
c.mu.Lock()
defer c.mu.Unlock()
return len(c.heard)
}
func eventually(t *testing.T, what string, is func() bool) {
t.Helper()
deadline := time.Now().Add(8 * time.Second)
for time.Now().Before(deadline) {
if is() {
return
}
time.Sleep(20 * time.Millisecond)
}
t.Fatalf("%s did not happen within the wait", what)
}
// A report published by a node reaches the controller, is recorded, and is acknowledged — so the
// stream does not hold it. A work queue is the check: what is acknowledged leaves it.
func TestNatsAReportIsHeardAndLeavesTheStream(t *testing.T) {
js := aBus(t)
store := &counted{}
_, stop := servingOn(t, js, store)
defer stop()
body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
eventually(t, "a report being recorded", func() bool { return store.count() == 1 })
eventually(t, "the report leaving the work queue", func() bool {
info, err := js.Context().StreamInfo("CONTROL")
return err == nil && info.State.Msgs == 0
})
}
// **The store window, in the server.** A report the store cannot take is naked with a delay and
// comes back; once the store is there it is recorded and leaves the stream. The controller holds
// nothing in the meantime — which is what the sequence check below is for: the message is still on
// the server while it waits.
func TestNatsAReportTheStoreCannotTakeIsHeldByTheServerAndComesBack(t *testing.T) {
js := aBus(t)
store := &counted{}
store.refusing(errors.Join(ErrTryAgain, errors.New("the database system is starting up")))
_, stop := servingOn(t, js, store)
defer stop()
body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
// Held: the message is the server's, unacknowledged, and still in the stream.
eventually(t, "the report being redelivered at least once", func() bool {
info, err := js.Context().ConsumerInfo("CONTROL", broker.ControllerName)
return err == nil && info.NumRedelivered >= 1
})
info, err := js.Context().StreamInfo("CONTROL")
if err != nil || info.State.Msgs != 1 {
t.Fatalf("a held report did not stay on the server: %+v, %v", info, err)
}
if store.count() != 0 {
t.Fatalf("a report was recorded by a store that was refusing it")
}
store.refusing(nil)
eventually(t, "the report being recorded once the store was back",
func() bool { return store.count() == 1 })
eventually(t, "the recorded report leaving the work queue", func() bool {
info, err := js.Context().StreamInfo("CONTROL")
return err == nil && info.State.Msgs == 0
})
}
// A report about a declaration the mesh has moved past is settled without being acted on, and
// leaves the stream rather than coming back for ever.
func TestNatsASupersededReportIsSettledAndNotActedOn(t *testing.T) {
js := aBus(t)
store := &sentAndHeardSafely{sent: "d2"}
_, stop := servingOn(t, js, store)
defer stop()
body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
eventually(t, "the superseded report leaving the stream", func() bool {
info, err := js.Context().StreamInfo("CONTROL")
return err == nil && info.State.Msgs == 0
})
if store.count() != 0 {
t.Fatalf("a report about a superseded declaration was acted on")
}
}
// sentAndHeardSafely is sentAndHeard, read from two goroutines.
type sentAndHeardSafely struct {
mu sync.Mutex
sent string
heard []Report
}
func (s *sentAndHeardSafely) Heard(_ context.Context, r Report) error {
s.mu.Lock()
defer s.mu.Unlock()
s.heard = append(s.heard, r)
return nil
}
func (s *sentAndHeardSafely) Outstanding(context.Context, string) (string, error) {
return s.sent, nil
}
func (s *sentAndHeardSafely) count() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.heard)
}
// **An enrolment answered through a reply address the stream would have eaten.**
//
// The caller waits on its own inbox and states that address in the request's payload. The check is
// that the answer arrives there — which is the whole reason the address is a field rather than the
// transport's reply, and this is the test design 25 §2 asks for so the reason cannot quietly become
// folklore.
func TestNatsAnEnrolmentIsAnsweredOnTheAddressInItsPayload(t *testing.T) {
js := aBus(t)
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()},
enroller: enrolsAs{reply: EnrolReply{Accepted: true, Node: "anchor"}}, log: quiet()}
ctx, stop := context.WithCancel(context.Background())
defer stop()
go func() { _ = s.Serve(ctx) }()
inbox := nats.NewInbox()
answers, err := js.Conn().SubscribeSync(inbox)
if err != nil {
t.Fatal(err)
}
body, _ := json.Marshal(EnrolRequest{Node: "anchor", Secret: "t", ReplyTo: inbox})
if _, err := js.Context().Publish(EnrolSubject, body); err != nil {
t.Fatal(err)
}
msg, err := answers.NextMsg(8 * time.Second)
if err != nil {
t.Fatalf("no answer reached the address the request named: %v", err)
}
var reply EnrolReply
if err := json.Unmarshal(msg.Data, &reply); err != nil {
t.Fatal(err)
}
if !reply.Accepted || reply.Node != "anchor" {
t.Fatalf("the answer was not the mesh's: %+v", reply)
}
// And the address really is not the one the transport carried: what the consumer saw was its
// own ack subject, which is why this had to travel in the payload.
if msg.Subject != inbox {
t.Fatalf("the answer arrived on %s, not the address the request named", msg.Subject)
}
}
// enrolsAs answers every request the same way.
type enrolsAs struct{ reply EnrolReply }
func (e enrolsAs) Enrol(context.Context, EnrolRequest) (EnrolReply, error) { return e.reply, nil }
// A heartbeat is core NATS: it reaches the controller and nothing is persisted, so the stream the
// reports live in stays empty.
func TestNatsAHeartbeatIsHeardAndNothingIsKept(t *testing.T) {
js := aBus(t)
store := &counted{}
_, stop := servingOn(t, js, store)
defer stop()
// Given time to subscribe: a core subscription that is not yet up misses what is published,
// which is the guarantee a heartbeat has and not a fault.
eventually(t, "the heartbeat subscription coming up", func() bool {
body, _ := json.Marshal(Alive{Node: "anchor"})
_ = js.Conn().Publish(AliveSubject("anchor"), body)
_ = js.Conn().Flush()
return store.count() >= 1
})
info, err := js.Context().StreamInfo("CONTROL")
if err != nil || info.State.Msgs != 0 {
t.Fatalf("a heartbeat was persisted, and the mesh's least valuable message now competes "+
"for retention with its most valuable: %+v, %v", info, err)
}
}
// The two events the controller follows arrive over one durable consumer with two filters, and it
// can acknowledge them.
//
// **Both halves are the point.** A consumer with several filter subjects is a 2.10 feature and this
// is the first thing in the mesh to use one; and a delivery from the events stream is acknowledged
// on a different ack subject from a delivery from the control stream, which the controller's own
// permission list has to cover or every announcement is redelivered for ever.
func TestNatsTheEventsTheControllerFollowsArriveAndAreAcknowledged(t *testing.T) {
js := aBus(t)
told := &toldAbout{}
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()}, log: quiet()}
if err := s.Follows(told); err != nil {
t.Fatal(err)
}
if err := s.Answers(replaysWith{}); err != nil {
t.Fatal(err)
}
ctx, stop := context.WithCancel(context.Background())
defer stop()
go func() { _ = s.Serve(ctx) }()
moved, _ := json.Marshal(Upgraded{Module: "gitea", Commit: "abcdef0123"})
if _, err := js.Context().Publish(broker.ControllerFollows[0], moved); err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(broker.ControllerFollows[1], []byte(`{}`)); err != nil {
t.Fatal(err)
}
eventually(t, "the catalogue's upgrade reaching the controller",
func() bool { return told.count() == 1 })
eventually(t, "both announcements being acknowledged", func() bool {
info, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName)
return err == nil && info.NumAckPending == 0 && info.Delivered.Consumer == 2
})
}
type toldAbout struct {
mu sync.Mutex
saw []Upgraded
fail error
}
func (u *toldAbout) Upgraded(_ context.Context, m Upgraded) error {
u.mu.Lock()
defer u.mu.Unlock()
if u.fail != nil {
return u.fail
}
u.saw = append(u.saw, m)
return nil
}
func (u *toldAbout) count() int {
u.mu.Lock()
defer u.mu.Unlock()
return len(u.saw)
}
// A store that never comes back: the report is let go once the bound passes, and it leaves the
// stream rather than being held for ever. The bound is the controller's, not the server's — nothing
// here sets max-deliver, and that is deliberate (streams.go).
func TestNatsAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
js := aBus(t)
store := &counted{}
store.refusing(errors.Join(ErrTryAgain, errors.New("connection refused")))
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()},
listener: store, log: quiet(), giveUp: 1500 * time.Millisecond}
ctx, stop := context.WithCancel(context.Background())
defer stop()
go func() { _ = s.Serve(ctx) }()
body, _ := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
eventually(t, "the report being let go once the bound passed", func() bool {
info, err := js.Context().StreamInfo("CONTROL")
return err == nil && info.State.Msgs == 0
})
if store.count() != 0 {
t.Fatalf("a report was recorded by a store that never came back")
}
}