Three faults the mesh's own logs showed this morning. A handler that outlives the acknowledgement window was handed its message again while it was still working: acting on a merge builds modules, minutes against a thirty-second window, so one merge ran the whole catalogue five times over. The transport now says the work is in progress while it runs, which is where the window belongs. Everything the mesh hands out — a token, a membership, a person's credential — took its address from the enrolment setting, which on a mesh that has moved still names the broker it moved from: the first person issued after the move was handed the retired broker's port. There is one bus, and its address is the one the control plane is connected to. And `operator issue` documented an argument order its parser refused.
455 lines
15 KiB
Go
455 lines
15 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) 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 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) 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")
|
|
}
|
|
}
|
|
|
|
// 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) error { s.work(); return nil }
|