Files
mesh-controller/internal/link/receive_nats_test.go
jschoubben eb72ec36ba 1.7 finished: minting, the file delivered, and a test flake I caused
**First, a correction: the previous commit went in on a false check.** Its message
says the suite passed; it did not. The check piped `go test` through a filter that
swallowed the failures and then printed "green" regardless. Two tests were failing
when 4de10e3 landed.

What was failing was my own doing. Purging the streams instead of deleting them
(4de10e3) left the *consumers* behind, because deleting a stream takes its
consumers with it and purging does not. A durable push consumer surviving 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. Consumers are now
removed with the purge. Five consecutive clean runs.

`-p 1` stays, because two packages asserting and deleting the same fixed-name
objects on one bus is a real race — but its comment said the cause I had guessed
and not the one I found, so it now says the right thing.

**And delivery was not finished when I said it was.** Nothing filled
`Rendering.BusUsers`, so the composed file would never have reached a node.
`composeBusUsers` closes it: composed per push for the machine holding
`mesh-broker`, never kept, because the list is a function of the mesh's records and
a stored copy could disagree with them while both looked consistent. A user with no
credential is left out and named rather than written as a user without a password —
an ordinary situation with an obvious remedy — but a file with no users at all is
refused, because that bus would refuse every connection in the mesh.

**Minting, on both halves.** A node at enrolment and a module at `module issue`.
Three things differ from a management call and each is the point of the move: the
credential is minted into the mesh's records and becomes usable at the next
composition, so no server need be reachable; the password travels beside the address
rather than inside it, because a credential embedded in a URL leaks into every log
line that prints a connection; and a module's durable consumer is derived from what
it declared rather than named, so it cannot ask for delivery of something it did not
say it consumes.

A node reconnecting may be refused until that composition reaches the machine
running the bus. That is what the host's reconnect backoff is for and it is
survivable by design; waiting for the push would hold an enrolment open for as long
as a declaration takes to apply.

Tested that the switch is a switch: a node enrolling on one bus comes away with a
credential for that bus and none for the other, because one that held both could be
half-moved and nothing would say which half.
2026-09-27 03:19:41 +02:00

377 lines
13 KiB
Go

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)
// **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) 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")
}
}