diff --git a/cmd/mesh-controller/ask.go b/cmd/mesh-controller/ask.go index 001d5e2..42423d5 100644 --- a/cmd/mesh-controller/ask.go +++ b/cmd/mesh-controller/ask.go @@ -37,7 +37,7 @@ func askCommand(ctx context.Context, args []string) error { arguments = json.RawMessage(positionals[2]) } - server, err := link.Connect(nil, nil) + server, err := connectLink(ctx, nil, nil, nil) if err != nil { return err } diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 600d36a..0a12c33 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -368,7 +368,7 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti } defer ident.Close() - server, err := link.Connect(nil, nil) + server, err := connectLink(ctx, nil, nil, nil) if err != nil { return err } @@ -472,7 +472,7 @@ func buildAndShow(ctx context.Context, source buildSource, path, ref string, wai return err } defer ident.Close() - server, err := link.Connect(nil, nil) + server, err := connectLink(ctx, nil, nil, nil) if err != nil { return err } diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 5ad9161..a4546e0 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -38,6 +38,33 @@ func reportUnhostable(node string, plan catalogue.Resolution) { // nothing in it was wrong, and no one edit was the one that should have been a new file. // serve is the control plane running: one connection to the broker, one queue, one consumer. +// connectLink opens the controller's link over whichever bus this process is on (design 25: one +// variable moves it). The streams and this controller's consumers are raised first on the new bus, +// so nothing served here finds them missing. +func connectLink(ctx context.Context, inv *inventory.Inventory, enroller link.Enroller, listener link.Listener) (*link.Server, error) { + busAddress, onNATS, err := broker.OnNATS() + if err != nil { + return nil, err + } + if err := broker.MustBeOneBus(os.Getenv(broker.AMQPVarName), busAddress); err != nil { + return nil, err + } + if !onNATS { + return link.Connect(enroller, listener) + } + if inv != nil { + if err := raiseTheBus(ctx, inv, busAddress); err != nil { + return nil, err + } + } + js, err := broker.Dial(busAddress) + if err != nil { + return nil, fmt.Errorf("the mesh is on the bus at %s and this control plane cannot reach it: %w", + busAddress, err) + } + return link.ConnectNats(js, enroller, listener), nil +} + func serve(ctx context.Context) error { open, err := openStores(ctx) if err != nil { @@ -90,7 +117,7 @@ func serve(ctx context.Context) error { work := link.Enrolment{Inventory: inv, Identity: ident, Management: management, Broker: known, OnNATS: onNATS} - server, err := link.Connect(work, work) + server, err := connectLink(ctx, inv, work, work) if err != nil { return err } @@ -100,11 +127,6 @@ func serve(ctx context.Context) error { // somebody deleted, a mesh raised from a restored backup, or a bus whose data directory was // replaced all have records and no objects — and a node whose consumer is missing hears nothing // while everything else about it looks correct. - if onNATS { - if err := raiseTheBus(ctx, inv, busAddress); err != nil { - return err - } - } // And build results nobody was waiting for. A build triggered any other way than `build` // would otherwise be reported into the void, which is the same as not reporting it. server.Records(builds{inv}) @@ -158,13 +180,13 @@ func declare(ctx context.Context, args []string) error { return err } - server, err := link.Connect(nil, nil) + server, err := connectLink(ctx, nil, nil, nil) if err != nil { return err } defer server.Close() - if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, node, raw, 15*time.Second); err != nil { + if err := link.Declare(ctx, server.Bus(), ident, node, raw, 15*time.Second); err != nil { return err } fmt.Printf("sent %s a signed declaration (%d bytes)\n", node, len(raw)) @@ -271,7 +293,7 @@ func pushCommand(ctx context.Context, args []string) error { return err } - server, err := link.Connect(nil, nil) + server, err := connectLink(ctx, nil, nil, nil) if err != nil { return err } @@ -332,7 +354,7 @@ func pushCommand(ctx context.Context, args []string) error { if err != nil { return err } - if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { + if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil { return err } // After it is away, not before. A digest recorded for something that failed to send would @@ -415,7 +437,7 @@ func pushCommand(ctx context.Context, args []string) error { return declarationWith(held, open, node, plan, settings, gens, Allocating) }, func(s readyNode, body []byte) error { - if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, + if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil { return err } @@ -628,7 +650,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error { len(refusals), strings.Join(refusals, "\n\n")) } - server, err := link.Connect(nil, nil) + server, err := connectLink(ctx, nil, nil, nil) if err != nil { return err } @@ -639,7 +661,7 @@ func sendTo(ctx context.Context, open *stores, names []string) error { if err != nil { return err } - if err := link.Declare(ctx, link.OverCurrent{Channel: server.Channel()}, ident, s.node, body, 15*time.Second); err != nil { + if err := link.Declare(ctx, server.Bus(), ident, s.node, body, 15*time.Second); err != nil { return err } record, err := inv.NodeByName(ctx, s.node) diff --git a/internal/link/bus_nats.go b/internal/link/bus_nats.go new file mode 100644 index 0000000..f4e4511 --- /dev/null +++ b/internal/link/bus_nats.go @@ -0,0 +1,73 @@ +package link + +import ( + "context" + "fmt" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-controller/internal/broker" +) + +// OverNats is the controller's outbound on the bus being built: the same three acts the other +// transport has, on the subjects the permissions were derived for (design 25). A declaration is a +// JetStream publish into NODES, where the node's own consumer waits for it; an event is announced on +// the subject its name derives to; a tool is asked by request and reply on the module's tool subject. +type OverNats struct{ JS *broker.JetStream } + +// declareSubject is where one node's declaration lands — the NODES stream's subject for it, and the +// only subject that node's consumer delivers. The host subscribes exactly this. +func declareSubject(node string) string { return "mesh.node." + node + ".declare" } + +func (b OverNats) PublishDeclaration(ctx context.Context, node string, body []byte) error { + publish, cancel := context.WithTimeout(ctx, 15*time.Second) + defer cancel() + if _, err := b.JS.Context().Publish(declareSubject(node), body, nats.Context(publish)); err != nil { + return fmt.Errorf("declaring to %s: %w", node, err) + } + return nil +} + +func (b OverNats) PublishEvent(ctx context.Context, key, source, node string, body []byte) error { + // The key is the subject: the controller's own events are named in full, and what a module + // emits is derived before it reaches here. Headers carry the envelope the other transport put + // in message properties (ADR 0042), so a consumer reads who and when without the payload. + msg := nats.NewMsg(key) + msg.Data = body + msg.Header.Set("x-source", source) + msg.Header.Set("x-node", node) + msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano)) + if err := b.JS.Conn().PublishMsg(msg); err != nil { + return fmt.Errorf("announcing %s: %w", key, err) + } + return nil +} + +func (b OverNats) AskTool(ctx context.Context, module, tool string, args []byte, timeout time.Duration) ([]byte, error) { + ask, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + reply, err := b.JS.Conn().RequestWithContext(ask, "mesh.mod."+module+".tool."+tool, args) + if err != nil { + return nil, fmt.Errorf("asking %s.%s: %w", module, tool, err) + } + return reply.Data, nil +} + +// ConnectNats is Connect for the bus being built: the controller's inbound and outbound over one +// JetStream connection the caller has already raised the streams on. Nothing is declared here — +// 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}, + js: js, + enroller: enroller, + listener: listener, + log: newLog(), + } +} + +// Bus is the controller's outbound, whichever transport it connected over. Callers that send a +// declaration or ask a tool use this rather than the channel, which one transport does not have. +func (s *Server) Bus() Bus { return s.bus } diff --git a/internal/link/serve.go b/internal/link/serve.go index f1c5750..c756754 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -7,6 +7,7 @@ import ( "encoding/json" "errors" "fmt" + "github.com/novox/mesh-controller/internal/broker" "log" "os" "time" @@ -70,6 +71,7 @@ type Server struct { bus Bus conn *amqp.Connection channel *amqp.Channel + js *broker.JetStream enroller Enroller listener Listener @@ -180,10 +182,12 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) { channel: channel, enroller: enroller, listener: listener, - log: log.New(os.Stdout, "", log.LstdFlags), + log: newLog(), }, nil } +func newLog() *log.Logger { return log.New(os.Stdout, "", log.LstdFlags) } + // Channel is the controller's channel, for the command line's own publishing. func (s *Server) Channel() *amqp.Channel { return s.channel } @@ -197,6 +201,9 @@ func (s *Server) Close() { if s.conn != nil { _ = s.conn.Close() } + if s.js != nil { + s.js.Close() + } } // Serve acts on what arrives until the context ends.