diff --git a/internal/broker/derived.go b/internal/broker/derived.go index e1da042..861b12c 100644 --- a/internal/broker/derived.go +++ b/internal/broker/derived.go @@ -28,6 +28,13 @@ type Consumer struct { // Queue is the queue group, set for a seat's worker so that "exactly one holder" survives a // seat later being relaxed to several. Authority and delivery are kept separate on purpose. Queue string + // Push asks the server to deliver to a subject rather than wait to be pulled. + // + // For the mesh's own consumer, where the controller wants every message to arrive in the one + // loop it already runs: pulling would mean a second goroutine fetching batches and handing + // them over, and a loop that acts on one message at a time is the property the store window + // depends on. A queue group implies this, because a group has nothing to pull from. + Push bool // AckWaitSeconds before an unacknowledged delivery is redelivered. AckWaitSeconds int // MaxDeliver before the message is dead-lettered; zero for the mesh's default. diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index 59f9caf..1867e21 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -39,6 +39,14 @@ func Dial(url string, opts ...nats.Option) (*JetStream, error) { return &JetStream{conn: conn, js: js}, nil } +// Conn is the connection itself, for what the mesh keeps off JetStream on purpose — a heartbeat, +// a tool call — where a lost message is answered by the next one or by a timeout the caller +// already handles (design 25 §3). +func (j *JetStream) Conn() *nats.Conn { return j.conn } + +// Context is the JetStream handle, for subscribing to what the consumers above define. +func (j *JetStream) Context() nats.JetStreamContext { return j.js } + func (j *JetStream) Close() { if j.conn != nil { j.conn.Close() @@ -112,7 +120,7 @@ func (j *JetStream) EnsureConsumer(c Consumer) error { // A queue group needs a delivery subject: a pull consumer has no group, and declaring one // without the other is refused by the server with a message that does not say which half is // missing. - if c.Queue != "" { + if c.Queue != "" || c.Push { want.DeliverSubject = "_DELIVER." + c.Name } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 80089fd..1dbb11f 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -143,6 +143,16 @@ func PermissionsFor(p Principal) (Permissions, error) { pub = []string{"mesh.control.>", "mesh.node.>", "mesh.build.>", "$JS.API.>"} sub = []string{"mesh.control.>", "mesh.build.>", "$JS.API.>"} + // The two events it reacts to, and its ack subject on the stream they arrive from + // (streams.go). **Each named, not a pattern**: `mesh.mod.*.event.>` would make the + // controller a subscriber to every event in the mesh, and its permission list would stop + // saying what it is for. The ack grant below is scoped per stream because the controller's + // consumer name is the same on both and `$JS.ACK.CONTROL.controller.>` does not cover a + // delivery from EVENTS — a consumer that cannot ack has every message redelivered for + // ever, refused by the list it already has. + sub = append(sub, ControllerFollows...) + pub = append(pub, "$JS.ACK.EVENTS."+ControllerName+".>") + case KindPerson: // Tools, and nothing else. Every subject a person may publish is a tool call; a person // who could publish an event would be able to claim a module said something. diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 90302ff..7b0d11a 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -150,3 +150,79 @@ func Overlaps() []string { sort.Strings(clashes) return clashes } + +// The mesh's own consumers. +// +// A seat's streams and a module's consumers are derived from declarations (derived.go). These two +// are not: **the controller is not a module and files no manifest**, so its authority and its +// subscriptions cannot come from a declaration that does not exist. They are named here, where the +// mesh's own streams are named, and narrowly — a controller subscribing `mesh.mod.*.event.>` would +// hear every event in the mesh, which it has no business doing and which would make its permission +// list stop explaining anything. + +// ControllerName is the controller's durable consumer on each stream it reads, and the name its +// ack subject is derived from (nats.go: `$JS.ACK..controller.>`). +const ControllerName = "controller" + +// ControllerFollows are the events the controller reacts to: the catalogue saying a module's +// current version moved, and a catalogue that has just started saying it may have missed builds. +// +// **These carry the local names the manifests hold today**, which still spell an event the way a +// routing key on the bus the mesh has does — `module..` rather than design 29's bare +// verb — so the derived subject names the module twice. It is consistent, and it is what the +// catalogue actually publishes, so it is what the controller must listen to. It changes when those +// names are converted, and not before: a subscription written against the name design 29 specifies +// would be a controller listening to a subject nothing publishes. +var ControllerFollows = []string{ + "mesh.mod.mesh-catalog.event.module.mesh-catalog.upgraded", + "mesh.mod.mesh-catalog.event.module.mesh-catalog.catching-up", +} + +// MeshConsumers is what the controller consumes, in the order a person reads it. +// +// **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's, +// not the server's (window.go): a message is held with a nak-and-delay until the controller either +// takes it or gives up and says so. A max-deliver here would dead-letter a push that was being +// held through a store restart — the exact message the stream exists to protect — some minutes +// before the controller had finished deciding about it. +func MeshConsumers() []Consumer { + return []Consumer{ + { + Name: ControllerName, + Stream: "CONTROL", + Push: true, + AckWaitSeconds: 30, + Why: "the controller is the single consumer of what nodes say; explicit ack and no " + + "max-deliver, because the store window's bound is the controller's own", + }, + { + Name: ControllerName, + Stream: "EVENTS", + Filters: ControllerFollows, + Push: true, + AckWaitSeconds: 30, + MaxDeliver: 5, + Why: "the two events the mesh's own controller reacts to; after max-deliver it " + + "dead-letters, because an announcement it cannot act on will not become actionable", + }, + } +} + +// Ensurer is the part of a JetStream connection consumer assertion needs, narrow for the reason +// Asserter is. +type Ensurer interface { + EnsureConsumer(c Consumer) error +} + +// AssertMeshConsumers brings the controller's own consumers into being, and says which one failed. +// +// After the streams, necessarily: a consumer on a stream that does not exist is refused, and the +// refusal names the stream rather than the order. +func AssertMeshConsumers(e Ensurer) error { + for _, c := range MeshConsumers() { + if err := e.EnsureConsumer(c); err != nil { + return fmt.Errorf("asserting consumer %s on %s: %w", c.Name, c.Stream, err) + } + } + return nil +} diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index c879021..edfd456 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -20,8 +20,8 @@ accounts { MESH { users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.API.>", "mesh.build.>", "mesh.control.>", "mesh.node.>"] } - subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.build.>", "mesh.control.>"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "mesh.build.>", "mesh.control.>", "mesh.node.>"] } + subscribe: { allow: ["$JS.API.>", "_INBOX.controller.>", "mesh.build.>", "mesh.control.>", "mesh.mod.mesh-catalog.event.module.mesh-catalog.catching-up", "mesh.mod.mesh-catalog.event.module.mesh-catalog.upgraded"] } allow_responses: { max: 1, ttl: "1m" } } } { user: "enrolment", password: "$2a$11$eeeeeeeeeeeeeeeeeeeeee", permissions: { diff --git a/internal/link/bus.go b/internal/link/bus.go index d35b864..b96d3cf 100644 --- a/internal/link/bus.go +++ b/internal/link/bus.go @@ -92,6 +92,34 @@ type OverNATS struct { JS nats.JetStreamContext } +// The subjects a node publishes on, and the controller listens to. +// +// One tree, and each name says who it is about: `mesh.control..…` is a node's own, which is +// what lets a node's account be granted exactly its own prefix and nothing of any other node's +// (design 25 §2, §4). The two that belong to no node — an enrolment, because a machine enrolling +// has no name the mesh has agreed to yet, and a build's outcome, because a builder is not +// reporting about itself — are named directly. +const ( + // EnrolSubject is where a joining machine asks. Its enrolment user may publish here and + // nowhere else, so a leaked token buys nothing but the chance to enrol. + EnrolSubject = "mesh.control.enrol" + + // BuiltSubject is where a build's outcome lands, for results nobody was waiting for. + BuiltSubject = "mesh.control.built" + + // AliveSubjects is every node's heartbeat. Core NATS, never a stream: a lost heartbeat is the + // next heartbeat, and a stream of them is the mesh's least valuable message competing for + // retention with its most valuable (design 25 §3). + AliveSubjects = "mesh.control.*.alive" +) + +// ReportSubject is where one node says what it did. On the CONTROL stream, because it is the +// message the store-window guarantee is about (ADR 0083). +func ReportSubject(node string) string { return "mesh.control." + node + ".report" } + +// AliveSubject is one node's heartbeat. +func AliveSubject(node string) string { return "mesh.control." + node + ".alive" } + // EventSubject is where a module's event lands. Derived from the emitter, never taken from the // caller: a source that could differ from the subject is an envelope that can lie about its // origin, and on NATS the account's permissions make the subject the authority (design 29 §2). diff --git a/internal/link/protocol.go b/internal/link/protocol.go index d814d15..df7276e 100644 --- a/internal/link/protocol.go +++ b/internal/link/protocol.go @@ -80,6 +80,20 @@ type EnrolRequest struct { // it. Nil from a node that found none, which is every converged one. Tunnel *Tunnel `json:"tunnel,omitempty"` + // ReplyTo is where the answer goes, as a field of the request rather than the transport's own + // reply address. + // + // **Because a stream eats the transport's field** (design 25 §2, verified against a running + // server): a message a JetStream consumer delivers has had its reply field claimed for that + // consumer's own ack address, so by the time the controller sees an enrolment, the field names + // where the *controller* must acknowledge, not where the node is waiting. An enrolment is the + // case that matters — a caller waiting on an ephemeral inbox, over a subject the store window + // may legitimately delay by several nak cycles. + // + // Empty on the bus the mesh runs on today, where the delivery carries the reply queue and the + // field means what it has always meant. + ReplyTo string `json:"reply_to,omitempty"` + // Redelivered is set by the control plane, never sent: the broker handed this request over a // second time. Such a request does not finish an enrolment already spent — the first time may // have answered, and the node holds what it was told. diff --git a/internal/link/receive_nats.go b/internal/link/receive_nats.go new file mode 100644 index 0000000..3acac49 --- /dev/null +++ b/internal/link/receive_nats.go @@ -0,0 +1,285 @@ +package link + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// The consume side on the bus being built. +// +// The shape is the AMQP one's, because the seam made them comparable: one loop, one message at a +// time, and the same window deciding. What differs is where a held message lives — and that is the +// whole point of the move. On the bus the mesh has, holding one means keeping an unacknowledged +// delivery in this process, bounded by the prefetch and lost if the controller stops. Here it is a +// `nak` with a delay: the message stays the server's, the controller keeps nothing but the moment +// it first could not take it, and a controller that restarts mid-window has nothing to lose. + +// natsInbound consumes what nodes and modules say over NATS. +type natsInbound struct { + js *broker.JetStream + // follows is the kinds asked for beyond what nodes say (Also). The events those are are the + // only ones the controller subscribes, and only when something is listening. + follows map[string]bool + // since is when the controller first could not take a message, by that message's place in its + // stream. + // + // **A timestamp, not a message.** This is the whole difference the move buys: the AMQP side + // keeps the delivery, and this keeps eight bytes saying when the window opened. A controller + // that restarts loses these and starts the window again, which is correct — it is holding + // nothing, and the messages are all still on the server. + since map[uint64]time.Time +} + +// Nats is the consume side of the bus being built. +func Nats(js *broker.JetStream) Inbound { + return &natsInbound{js: js, follows: map[string]bool{}, since: map[uint64]time.Time{}} +} + +// Also records one more kind to subscribe. Nothing is subscribed here: the controller's consumer on +// the events stream carries both of these as filters, so it is created once, in Receive, with +// whatever was asked for — and not at all when nothing was. +func (n *natsInbound) Also(kind string) error { + switch kind { + case KindModuleMoved, KindCatchUp: + n.follows[kind] = true + return nil + default: + return fmt.Errorf("nothing subscribes %s separately on this bus", kind) + } +} + +func (n *natsInbound) Close() {} + +// Receive consumes until the context ends. +// +// Three subscriptions, and each is a channel the one loop selects on. **Channels rather than +// callbacks**: the library would run a handler on its own goroutine, and the window's bookkeeping — +// which message is held, and since when — is read and written without a lock because the AMQP loop +// never had two. A second goroutine would make that wrong in a way no test would catch. +func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Control)) error { + if err := broker.AssertMeshConsumers(n.js); err != nil { + return err + } + js, conn := n.js.Context(), n.js.Conn() + + // What nodes say, off the CONTROL stream. Bound to the durable the controller asserted rather + // than creating one here: the consumer is an object with a configuration — ack policy, ack + // wait, redelivery — and a client that creates its own would be a second opinion about it. + control := make(chan *nats.Msg, Prefetch) + said, err := js.ChanSubscribe("", control, nats.Bind("CONTROL", broker.ControllerName)) + if err != nil { + return fmt.Errorf("subscribing to what nodes say: %w", err) + } + defer func() { _ = said.Unsubscribe() }() + + // Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because + // they are their own guarantee: a lost one is the next one. + beats := make(chan *nats.Msg, Prefetch) + alive, err := conn.ChanSubscribe(AliveSubjects, beats) + if err != nil { + return fmt.Errorf("subscribing to heartbeats: %w", err) + } + defer func() { _ = alive.Unsubscribe() }() + + // The events the controller follows, when something is listening for them. + var events chan *nats.Msg + if len(n.follows) > 0 { + events = make(chan *nats.Msg, Prefetch) + followed, err := js.ChanSubscribe("", events, nats.Bind("EVENTS", broker.ControllerName)) + if err != nil { + return fmt.Errorf("subscribing to what the catalogue says: %w", err) + } + defer func() { _ = followed.Unsubscribe() }() + } + + // A connection that dropped is said, not discovered. A controller whose bus connection is gone + // is a mesh where nothing can be told anything. + gone := make(chan error, 1) + conn.SetDisconnectErrHandler(func(_ *nats.Conn, err error) { + select { + case gone <- err: + default: + } + }) + + for { + select { + case <-ctx.Done(): + return nil + case err := <-gone: + return fmt.Errorf("the bus connection dropped: %w", err) + case msg := <-beats: + n.deliver(ctx, act, msg, false) + case msg := <-events: + n.deliver(ctx, act, msg, true) + case msg, ok := <-control: + if !ok { + return errors.New("the bus stopped delivering") + } + n.deliver(ctx, act, msg, true) + } + } +} + +// deliver names one message and hands it to the loop, or drops it where the mesh has no name for +// its subject — which cannot happen through a filter the controller wrote, and is said rather than +// ignored for exactly that reason. +func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Control), + msg *nats.Msg, streamed bool) { + + kind, known := kindOfSubject(msg.Subject) + if !known { + if streamed { + _ = msg.Term() + } + return + } + m := &natsControl{kind: kind, msg: msg, on: n} + if streamed { + // A message with no metadata is not from a stream, whatever it was delivered on, and the + // window has nothing to hold it by. Said by leaving the sequence at zero. + if meta, err := msg.Metadata(); err == nil { + m.seq = meta.Sequence.Stream + m.delivered = meta.NumDelivered + } + } + act(ctx, m) +} + +// kindOfSubject is how this transport's addressing becomes what the mesh calls a message. +// +// By subject, which is the only thing the server enforces: a body claiming to be a report does not +// make it one, and on this bus the subject an account may publish *is* its authority (design 29 +// §2). The mirror of the routing-key table on the bus the mesh has. +func kindOfSubject(subject string) (string, bool) { + switch subject { + case EnrolSubject: + return KindEnrolment, true + case BuiltSubject: + return KindBuilt, true + } + if node, rest, ok := strings.Cut(strings.TrimPrefix(subject, "mesh.control."), "."); ok && + node != "" && !strings.Contains(node, ".") { + switch rest { + case "report": + return KindReport, true + case "alive": + return KindHeartbeat, true + } + } + switch subject { + case broker.ControllerFollows[0]: + return KindModuleMoved, true + case broker.ControllerFollows[1]: + return KindCatchUp, true + } + return "", false +} + +// natsControl is one message from the bus being built, as the controller reads it. +type natsControl struct { + kind string + msg *nats.Msg + on *natsInbound + // seq is this message's place in its stream; zero for a core message, which has none and + // cannot be held. + seq uint64 + // delivered is how many times the server has handed this message over, this time included. + delivered uint64 +} + +func (m *natsControl) Kind() string { return m.kind } +func (m *natsControl) Body() []byte { return m.msg.Data } + +// Redelivered is what the server counted, not what the controller remembers. Which is the answer to +// a question the AMQP side could only guess at across a restart: an enrolment redelivered because +// the controller stopped mid-answer reads as redelivered to the controller that comes back. +func (m *natsControl) Redelivered() bool { return m.delivered > 1 } + +func (m *natsControl) HeldFor() time.Duration { + if m.seq == 0 { + return 0 + } + first, held := m.on.since[m.seq] + if !held { + return 0 + } + return time.Since(first) +} + +// About is nothing here, and that is the point. +// +// Setting a held message aside when a newer one about the same thing arrives is what a controller +// holding deliveries in memory can do. A naked message belongs to the server and comes back +// whatever happened meanwhile, so the question "is this the past?" is answered by what the message +// says instead — the digest of the declaration a report is about (window.go, design 25 §3). +func (m *natsControl) About(string) {} + +// Answer publishes to the reply subject the request carries **in its payload**. +// +// Not `Respond`, and not the message's reply field: a message a JetStream consumer delivers has had +// that field claimed for the consumer's own ack address, so answering it would send the reply to +// `$JS.ACK.CONTROL.controller.…` and the enrolling node would wait out its timeout. Verified +// against a running server (design 25 §2), which is why it is a field of the request and this reads +// it from there. +func (m *natsControl) Answer(ctx context.Context, body []byte) error { + var addressed replyAddressed + if err := json.Unmarshal(m.msg.Data, &addressed); err != nil { + return fmt.Errorf("that request cannot be read, so its reply address cannot be: %w", err) + } + if addressed.ReplyTo == "" { + return errors.New("that request named no reply subject in its payload, so nothing can be " + + "told the answer") + } + return m.on.js.Conn().PublishMsg(&nats.Msg{Subject: addressed.ReplyTo, Data: body}) +} + +func (m *natsControl) Took() error { + m.forget() + if m.seq == 0 { + // Core NATS: nothing is keeping it, so there is nothing to settle. + return nil + } + return m.msg.Ack(nats.Context(context.Background())) +} + +// Drop terminates the delivery: understood, and the server is told not to send it again. Different +// from an ack only in the server's own accounting, which is where somebody asking "what happened to +// that message" will look. +func (m *natsControl) Drop() error { + m.forget() + if m.seq == 0 { + return nil + } + return m.msg.Term() +} + +// Hold hands the message back with a delay, and remembers when the window opened. +func (m *natsControl) Hold(after time.Duration) error { + if m.seq == 0 { + return errors.New("a message that is not in a stream cannot be held: nothing is keeping it") + } + if _, already := m.on.since[m.seq]; !already { + m.on.since[m.seq] = time.Now() + } + return m.msg.NakWithDelay(after) +} + +func (m *natsControl) forget() { + if m.on != nil && m.seq != 0 { + delete(m.on.since, m.seq) + } +} + +// replyAddressed is the one field every message that expects an answer carries. +type replyAddressed struct { + ReplyTo string `json:"reply_to,omitempty"` +} diff --git a/internal/link/receive_nats_test.go b/internal/link/receive_nats_test.go new file mode 100644 index 0000000..2f4a58e --- /dev/null +++ b/internal/link/receive_nats_test.go @@ -0,0 +1,365 @@ +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) + + // The mesh's own streams and consumers, asserted the way the controller asserts them — and + // torn down after, so one test's held message is never another's surprise. + for _, s := range broker.MeshStreams() { + _ = js.Context().DeleteStream(s.Name) + } + if err := broker.AssertMeshStreams(js); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + for _, s := range broker.MeshStreams() { + _ = js.Context().DeleteStream(s.Name) + } + }) + 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") + } +}