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).
This commit is contained in:
@@ -9,12 +9,13 @@ import (
|
|||||||
// the streams and the controller's consumers are asserted by Raise, before anything is served.
|
// 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 {
|
func ConnectNats(js *broker.JetStream, enroller Enroller, listener Listener) *Server {
|
||||||
return &Server{
|
return &Server{
|
||||||
inbound: Nats(js),
|
inbound: Nats(js),
|
||||||
bus: OverNATS{JS: js.Context(), Conn: js.Conn()},
|
bus: OverNATS{JS: js.Context(), Conn: js.Conn()},
|
||||||
js: js,
|
js: js,
|
||||||
enroller: enroller,
|
consumers: js,
|
||||||
listener: listener,
|
enroller: enroller,
|
||||||
log: newLog(),
|
listener: listener,
|
||||||
|
log: newLog(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -73,6 +73,8 @@ type Server struct {
|
|||||||
inbound Inbound
|
inbound Inbound
|
||||||
bus Bus
|
bus Bus
|
||||||
js *broker.JetStream
|
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
|
enroller Enroller
|
||||||
listener Listener
|
listener Listener
|
||||||
@@ -383,6 +385,7 @@ func (s *Server) enrolling(ctx context.Context, m Control) {
|
|||||||
default:
|
default:
|
||||||
reply = accepted
|
reply = accepted
|
||||||
s.log.Printf("enrolled %s", accepted.Node)
|
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()
|
_ = 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.
|
// 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
|
// This is for results nobody was waiting for. A build asked for with `build` is answered directly
|
||||||
|
|||||||
Reference in New Issue
Block a user