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.
560 lines
19 KiB
Go
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 }
|