Files
mesh-controller/internal/link/receive_nats_test.go
T
jschoubben 1ebad3786c The mesh says what it applied, and the replay has an address it may use
The pipeline was observable from a merge to an artifact and went dark where it touched a machine: a
node's report is control traffic only the control plane reads, so nothing said which version a
machine runs, or that it refused to (novox/hq ADR 0134). The control plane now states both under the
seat it holds — a role's events belong to the role and keep their address when the holder is
replaced — and only when the report is news, because a machine reconciles every minute and a fact per
report would be a fact per minute per machine.

Whether a report is news is the store's answer: it holds the previous one, so the listener returns it
and the server states the fact. That also gives the catch-up replay a subject the controller may
publish: it was published as a module's event from a module called "control-plane", which does not
exist, so the controller's own account refused it and every catalogue that asked what it missed was
answered with nothing.
2026-09-28 16:07:18 +02:00

560 lines
19 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"io"
"log"
"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)
// **Streams purged, consumers removed.** Both halves, and each was learned by getting it wrong.
//
// The streams are emptied rather than deleted and recreated, because delete-then-add is not a
// reset: the server's teardown races the creation, and a test then inherits the previous one's
// messages — which reads as a redelivery bug in the code under test.
//
// The consumers are removed, because deleting a stream used to take them with it and purging
// does not. A durable *push* consumer that survives between tests keeps pushing to a delivery
// subject the previous test's subscription has gone from: the messages count as delivered, go
// nowhere, and the next test waits out its timeout for an announcement the server believes it
// already sent. The controller recreates what it needs on start, so leaving none is correct.
if err := broker.AssertMeshStreams(js); err != nil {
t.Fatal(err)
}
clean := func() {
for _, c := range broker.MeshConsumers() {
_ = js.Context().DeleteConsumer(c.Stream, c.Name)
}
for _, s := range broker.MeshStreams() {
_ = js.Context().PurgeStream(s.Name)
}
}
clean()
t.Cleanup(clean)
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) (bool, error) {
c.mu.Lock()
defer c.mu.Unlock()
if c.err != nil {
return false, c.err
}
c.heard = append(c.heard, r)
// News, so what the mesh states about a report is exercised wherever a report is.
return true, 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 quiet() *log.Logger { return log.New(io.Discard, "", 0) }
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) (bool, error) {
s.mu.Lock()
defer s.mu.Unlock()
s.heard = append(s.heard, r)
return true, 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")
}
}
// A merge announcement is not what these tests are about; taken and forgotten.
func (t *toldAbout) SourceMoved(context.Context, SourceMoved) error { return nil }
// Every subject the controller follows decodes to a kind, and decoding never reaches past the
// list: the day the list was three long and the decoder named a fourth, every message panicked
// the control plane (2026-09-28).
func TestEverySubjectTheControllerFollowsDecodesToAKind(t *testing.T) {
for _, subject := range broker.ControllerFollows {
if _, ok := kindOfSubject(subject); !ok {
t.Errorf("%s is followed and decodes to nothing", subject)
}
}
if _, ok := kindOfSubject("mesh.mod.nobody.event.nothing"); ok {
t.Error("a subject nobody follows decoded to a kind")
}
}
// **A handler slower than the acknowledgement window is not handed its message again.**
//
// The bus waits a fixed time to be told a message was taken and then redelivers, which is right for
// a consumer that died and wrong for one that is busy. Acting on a merge builds modules — minutes
// against a thirty-second window — and the same merge was handed over five times while the first
// build was still running (2026-09-28). Here the window is two seconds and the work takes six.
func TestNatsWorkSlowerThanTheWindowIsNotHandedOverAgain(t *testing.T) {
js := aBus(t)
// The controller's own consumer, with a window short enough to outlive in a test.
if err := js.EnsureConsumer(broker.Consumer{
Name: broker.ControllerName, Stream: "CONTROL", Push: true, AckWaitSeconds: 2,
Why: "a window short enough to outlive in a test",
}); err != nil {
t.Fatal(err)
}
var mu sync.Mutex
handled := 0
slow := make(chan struct{})
s, stop := servingOn(t, js, nil)
defer stop()
s.listener = slowly{func() {
mu.Lock()
handled++
first := handled == 1
mu.Unlock()
if first {
time.Sleep(6 * time.Second)
close(slow)
}
}}
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
select {
case <-slow:
case <-time.After(30 * time.Second):
t.Fatal("the slow work never finished")
}
// A moment for a redelivery to arrive, if the bus were going to send one.
time.Sleep(3 * time.Second)
mu.Lock()
defer mu.Unlock()
if handled != 1 {
t.Fatalf("one report was handled %d times, so slow work is run again while it is running", handled)
}
}
// slowly is a listener that runs whatever it was given.
type slowly struct{ work func() }
func (s slowly) Heard(context.Context, Report) (bool, error) { s.work(); return true, nil }
// **The mesh says what it applied** (novox/hq ADR 0134), under the seat the control plane holds — and
// says nothing when a report is the same state said again, which is what a machine reconciling every
// minute sends.
func TestNatsTheMeshSaysWhatAMachineApplied(t *testing.T) {
js := aBus(t)
heard := make(chan *nats.Msg, 4)
sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) {
heard <- m
})
if err != nil {
t.Fatal(err)
}
defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection
_, stop := servingOn(t, js, &counted{})
defer stop()
// A report that changed something: the store says it was news.
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store", "broker"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
select {
case m := <-heard:
if m.Subject != SeatEventSubject(MeshControllerSeat, KeyApplied) {
t.Fatalf("the mesh stated %q", m.Subject)
}
var said Applied
if err := json.Unmarshal(m.Data, &said); err != nil {
t.Fatal(err)
}
if said.Node != "anchor" || said.Declared != "d1" || said.Resources != 2 {
t.Fatalf("it said %+v", said)
}
case <-time.After(10 * time.Second):
t.Fatal("the mesh said nothing about a machine that now runs something else")
}
// A refusal is its own fact, with the reason in it rather than only in a log.
refusal, err := json.Marshal(Report{Node: "anchor", Declared: "d2",
Failed: map[string]string{"gitea.server": "no such image"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), refusal); err != nil {
t.Fatal(err)
}
select {
case m := <-heard:
if m.Subject != SeatEventSubject(MeshControllerSeat, KeyRefused) {
t.Fatalf("a refusal was stated as %q", m.Subject)
}
var said Refused
if err := json.Unmarshal(m.Data, &said); err != nil {
t.Fatal(err)
}
if said.Failed["gitea.server"] == "" {
t.Fatalf("the refusal does not say which resource or why: %+v", said)
}
case <-time.After(10 * time.Second):
t.Fatal("the mesh said nothing about a machine that refused what it was sent")
}
}
// And a report that is not news is not a fact. A machine reconciles every minute; a fact per report
// would be a fact per minute per machine, which is a stream nobody reads.
func TestNatsAReportThatIsNotNewsIsNotStated(t *testing.T) {
js := aBus(t)
heard := make(chan *nats.Msg, 4)
sub, err := js.Conn().Subscribe(SeatEventSubject(MeshControllerSeat, ">"), func(m *nats.Msg) {
heard <- m
})
if err != nil {
t.Fatal(err)
}
defer sub.Unsubscribe() //nolint:errcheck // the subscription dies with the connection
// A store that records the report and says it was nothing new — which is what the mesh's own
// store says about a machine repeating itself.
_, stop := servingOn(t, js, sameAgain{})
defer stop()
body, err := json.Marshal(Report{Node: "anchor", Declared: "d1", Applied: []string{"store"}})
if err != nil {
t.Fatal(err)
}
if _, err := js.Context().Publish(ReportSubject("anchor"), body); err != nil {
t.Fatal(err)
}
select {
case m := <-heard:
t.Fatalf("the mesh stated %q about a machine that changed nothing", m.Subject)
case <-time.After(3 * time.Second):
}
}
// sameAgain records a report and says it was the same state said again.
type sameAgain struct{}
func (sameAgain) Heard(context.Context, Report) (bool, error) { return false, nil }