From cfb6554e692b5bceff876104bd922131c5928e69 Mon Sep 17 00:00:00 2001 From: jochens Date: Fri, 2 Oct 2026 16:55:50 +0200 Subject: [PATCH] A node is given how it hears its declarations as it enrols The node's declaration consumer was asserted only when the control plane started, so the first machine of a mesh, enrolling after the control plane was up, joined and then heard nothing: its host retried 'consumer not found' for ever (novox/hq issue 146). --- internal/link/bus_nats.go | 13 +++---- internal/link/enrol_consumer_test.go | 51 ++++++++++++++++++++++++++++ internal/link/serve.go | 21 ++++++++++++ 3 files changed, 79 insertions(+), 6 deletions(-) create mode 100644 internal/link/enrol_consumer_test.go diff --git a/internal/link/bus_nats.go b/internal/link/bus_nats.go index b902da4..9b48ae9 100644 --- a/internal/link/bus_nats.go +++ b/internal/link/bus_nats.go @@ -9,12 +9,13 @@ import ( // the streams and the controller's consumers are asserted by Raise, before anything is served. func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server { return &Server{ - inbound: Nats(js), - bus: OverNATS{JS: js.Context(), Conn: js.Conn()}, - js: js, - enroller: enroller, - listener: listener, - log: newLog(), + inbound: Nats(js), + bus: OverNATS{JS: js.Context(), Conn: js.Conn()}, + js: js, + consumers: js, + enroller: enroller, + listener: listener, + log: newLog(), } } diff --git a/internal/link/enrol_consumer_test.go b/internal/link/enrol_consumer_test.go new file mode 100644 index 0000000..77a8c52 --- /dev/null +++ b/internal/link/enrol_consumer_test.go @@ -0,0 +1,51 @@ +package link + +import ( + "context" + "encoding/json" + "io" + "log" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/broker" +) + +type ensured struct{ consumers []broker.Consumer } + +func (e *ensured) EnsureConsumer(c broker.Consumer) error { + e.consumers = append(e.consumers, c) + return nil +} + +type acceptsAs string + +func (n acceptsAs) Enrol(context.Context, EnrolRequest) (EnrolReply, error) { + return EnrolReply{Accepted: true, Node: string(n)}, nil +} + +type anEnrolment struct{ body []byte } + +func (m anEnrolment) Kind() string { return "enrol" } +func (m anEnrolment) Body() []byte { return m.body } +func (m anEnrolment) Redelivered() bool { return false } +func (m anEnrolment) HeldFor() time.Duration { return 0 } +func (m anEnrolment) Answer(context.Context, []byte) error { return nil } +func (m anEnrolment) Took() error { return nil } +func (m anEnrolment) Hold(time.Duration) error { return nil } +func (m anEnrolment) Drop() error { return nil } + +// **A node that enrols can hear its declarations at once** (novox/hq 04-ISSUES/146): its consumer is +// made as it enrols, not only when the control plane next starts — the first machine of a mesh +// enrols after the control plane is up, and heard nothing. +func TestAnEnrolledNodeIsGivenHowItHearsItsDeclarations(t *testing.T) { + made := &ensured{} + s := &Server{enroller: acceptsAs("anchor"), consumers: made, log: log.New(io.Discard, "", 0)} + body, _ := json.Marshal(EnrolRequest{Node: "anchor"}) + s.enrolling(context.Background(), anEnrolment{body: body}) + + want := broker.NodeConsumer("anchor") + if len(made.consumers) != 1 || made.consumers[0].Name != want.Name || made.consumers[0].Stream != want.Stream { + t.Fatalf("the enrolled node was given %v, want its own declaration consumer %v", made.consumers, want) + } +} diff --git a/internal/link/serve.go b/internal/link/serve.go index 1099873..d6aea65 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -73,6 +73,8 @@ type Server struct { inbound Inbound bus Bus js *broker.JetStream + // consumers makes a node's declaration consumer as it enrols; the bus connection, or a stand-in. + consumers interface{ EnsureConsumer(broker.Consumer) error } enroller Enroller listener Listener @@ -383,6 +385,7 @@ func (s *Server) enrolling(ctx context.Context, m Control) { default: reply = accepted s.log.Printf("enrolled %s", accepted.Node) + s.hearsItsDeclarations(accepted.Node) } } @@ -402,6 +405,24 @@ func (s *Server) enrolling(ctx context.Context, m Control) { _ = m.Took() } +// hearsItsDeclarations makes the consumer a node reads its declarations through, as it enrols and +// before it is answered. +// +// **Created at enrolment, as the consumer's own doc has always said** (novox/hq 04-ISSUES/146). It +// was asserted only when the control plane started, so the first machine of a mesh — which enrols +// after the control plane is already up — joined and then heard nothing, its host retrying "consumer +// not found" for ever. Failing here is said and does not unspend the token: the next start of the +// control plane asserts it again. +func (s *Server) hearsItsDeclarations(node string) { + if s.consumers == nil { + return + } + if err := s.consumers.EnsureConsumer(broker.NodeConsumer(node)); err != nil { + s.log.Printf("%s enrolled, and how it hears its declarations could not be made — it will hear "+ + "nothing until the control plane next starts: %v", node, err) + } +} + // wasBuilt keeps what a builder said, whichever way it went. // // This is for results nobody was waiting for. A build asked for with `build` is answered directly