**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 when4de10e3landed. 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.
377 lines
13 KiB
Go
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")
|
|
}
|
|
}
|