The control queue was bound to enrol and not to report, so every report a node sent was accepted by the broker, matched no binding, and dropped. The publisher saw success and the consumer saw nothing, for an afternoon. The refactor that was meant to bind both never applied -- it left behind a helper nothing called, which compiled and passed vet. The loop is now where the bind is, so there is one place to forget rather than two.
224 lines
7.9 KiB
Go
224 lines
7.9 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"crypto/ed25519"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"os"
|
|
"strings"
|
|
"time"
|
|
|
|
amqp "github.com/rabbitmq/amqp091-go"
|
|
)
|
|
|
|
// AMQPVar is the control plane's own connection to the broker.
|
|
const AMQPVar = "MESH_BROKER_AMQP"
|
|
|
|
// Enroller is what the control plane does with an enrolment request.
|
|
//
|
|
// An interface so the serving loop can be tested against a real broker without a database, and
|
|
// so the two concerns — moving messages, and deciding — stay apart.
|
|
type Enroller interface {
|
|
// Enrol spends the token, records the key, and reports the node's name. The error is
|
|
// returned to the node as a refusal; it must be the same for every reason a token can fail.
|
|
Enrol(ctx context.Context, secret string, public ed25519.PublicKey, profile map[string]any) (EnrolReply, error)
|
|
}
|
|
|
|
// Server consumes what nodes say.
|
|
type Server struct {
|
|
conn *amqp.Connection
|
|
channel *amqp.Channel
|
|
enroller Enroller
|
|
log *log.Logger
|
|
}
|
|
|
|
// Connect opens the control plane's own connection to the broker.
|
|
func Connect(enroller Enroller) (*Server, error) {
|
|
url := strings.TrimSpace(os.Getenv(AMQPVar))
|
|
if url == "" {
|
|
return nil, fmt.Errorf(
|
|
"this control plane has no %s, so it cannot reach its broker. Nodes talk to it over "+
|
|
"the broker and nowhere else, so without this it can hold records and answer "+
|
|
"nothing", AMQPVar)
|
|
}
|
|
|
|
conn, err := amqp.Dial(url)
|
|
if err != nil {
|
|
// Not quoted back: the URL carries the control plane's own broker password.
|
|
return nil, fmt.Errorf("cannot reach the broker named in %s: %w", AMQPVar, err)
|
|
}
|
|
channel, err := conn.Channel()
|
|
if err != nil {
|
|
conn.Close()
|
|
return nil, err
|
|
}
|
|
|
|
// Declared here rather than assumed. The control plane is the only thing that may create
|
|
// them — a node's account can write to this exchange and read its own queue, and configure
|
|
// nothing else, so a node arriving before the control plane has ever run finds nothing and
|
|
// says so, rather than quietly creating a topology nobody designed.
|
|
if err := channel.ExchangeDeclare(Exchange, "direct", true, false, false, false, nil); err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("cannot declare the %s exchange: %w", Exchange, err)
|
|
}
|
|
if _, err := channel.QueueDeclare(ControlQueue, true, false, false, false, nil); err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("cannot declare the %s queue: %w", ControlQueue, err)
|
|
}
|
|
// Every key a node may publish. Binding one and forgetting another is a message the broker
|
|
// accepts, finds no queue for, and drops — the publisher sees success and the consumer sees
|
|
// nothing. That is exactly what happened to reports: `report` was left unbound while `enrol`
|
|
// worked, so nodes announced what they had applied into a void for an afternoon.
|
|
for _, key := range []string{KeyEnrol, KeyReport} {
|
|
if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err)
|
|
}
|
|
}
|
|
|
|
return &Server{conn: conn, channel: channel, enroller: enroller,
|
|
log: log.New(os.Stdout, "", log.LstdFlags)}, nil
|
|
}
|
|
|
|
// Channel is the control plane's channel, for sending declarations.
|
|
func (s *Server) Channel() *amqp.Channel { return s.channel }
|
|
|
|
func (s *Server) Close() {
|
|
if s.channel != nil {
|
|
_ = s.channel.Close()
|
|
}
|
|
if s.conn != nil {
|
|
_ = s.conn.Close()
|
|
}
|
|
}
|
|
|
|
// Serve consumes until the context ends.
|
|
//
|
|
// One consumer, deliberately: with two, the broker would round-robin between them and each would
|
|
// receive half of what it expects — a fault this project has already had, between a module's
|
|
// daemon and its capability server.
|
|
func (s *Server) Serve(ctx context.Context) error {
|
|
// Prefetch of one. The control plane writes to a database per message, and a burst of
|
|
// enrolments delivered all at once would be held in memory rather than left on the broker,
|
|
// which is the one place they survive a restart.
|
|
if err := s.channel.Qos(1, 0, false); err != nil {
|
|
return err
|
|
}
|
|
|
|
deliveries, err := s.channel.ConsumeWithContext(ctx, ControlQueue, "control-plane",
|
|
false, false, false, false, nil)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
closed := s.conn.NotifyClose(make(chan *amqp.Error, 1))
|
|
s.log.Printf("consuming %s, bound to %s/{%s,%s}", ControlQueue, Exchange, KeyEnrol, KeyReport)
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case reason := <-closed:
|
|
// Said rather than returned quietly. A control plane whose broker connection dropped
|
|
// is a mesh where nothing can be told anything, and the reason is the first thing
|
|
// anybody will want.
|
|
return fmt.Errorf("the broker connection closed: %v", reason)
|
|
case delivery, ok := <-deliveries:
|
|
if !ok {
|
|
return errors.New("the broker stopped delivering")
|
|
}
|
|
s.handle(ctx, delivery)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Server) handle(ctx context.Context, delivery amqp.Delivery) {
|
|
switch delivery.RoutingKey {
|
|
case KeyEnrol:
|
|
s.handleEnrol(ctx, delivery)
|
|
case KeyReport:
|
|
s.handleReport(delivery)
|
|
default:
|
|
// Rejected without requeue: a message nothing understands will not be understood on the
|
|
// next attempt either, and requeuing it would spin.
|
|
s.log.Printf("refusing a message with routing key %q", delivery.RoutingKey)
|
|
_ = delivery.Reject(false)
|
|
}
|
|
}
|
|
|
|
// handleReport records what a node says it did.
|
|
//
|
|
// A node states; nothing here writes anything the node claimed about itself beyond that it was
|
|
// heard from. What it applied is its own account of its own machine, and the mesh keeps the last
|
|
// one as a copy for recovery rather than as a source (novox/hq 09-the-node-lifecycle).
|
|
func (s *Server) handleReport(delivery amqp.Delivery) {
|
|
var report Report
|
|
if err := json.Unmarshal(delivery.Body, &report); err != nil {
|
|
s.log.Printf("a report could not be read: %v", err)
|
|
_ = delivery.Reject(false)
|
|
return
|
|
}
|
|
switch {
|
|
case report.Refused != "":
|
|
s.log.Printf("%s refused a declaration: %s", report.Node, report.Refused)
|
|
case len(report.Failed) > 0:
|
|
s.log.Printf("%s applied %d and failed: %v", report.Node, len(report.Applied), report.Failed)
|
|
default:
|
|
s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied))
|
|
}
|
|
_ = delivery.Ack(false)
|
|
}
|
|
|
|
func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) {
|
|
reply := EnrolReply{Refusal: "that token cannot be used"}
|
|
|
|
var request EnrolRequest
|
|
if err := json.Unmarshal(delivery.Body, &request); err != nil {
|
|
s.log.Printf("an enrolment request could not be read: %v", err)
|
|
} else {
|
|
accepted, err := s.enroller.Enrol(ctx, request.Secret, request.PublicKey, request.Profile)
|
|
if err != nil {
|
|
// Logged in full here, where an operator can see it; sent back as one refusal, so
|
|
// that somebody guessing learns nothing from which reason came back.
|
|
s.log.Printf("refusing enrolment for %q: %v", request.Node, err)
|
|
} else {
|
|
reply = accepted
|
|
s.log.Printf("enrolled %s", accepted.Node)
|
|
}
|
|
}
|
|
|
|
s.reply(ctx, delivery, reply)
|
|
|
|
// Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the
|
|
// request on the broker rather than having consumed it silently. Enrolment is idempotent
|
|
// only in the sense that the token is spent — a redelivery gets the refusal, which is
|
|
// correct and visible, where a lost request is neither.
|
|
_ = delivery.Ack(false)
|
|
}
|
|
|
|
func (s *Server) reply(ctx context.Context, delivery amqp.Delivery, reply EnrolReply) {
|
|
if delivery.ReplyTo == "" {
|
|
s.log.Print("an enrolment request named no reply queue, so nothing can be told the answer")
|
|
return
|
|
}
|
|
body, err := json.Marshal(reply)
|
|
if err != nil {
|
|
s.log.Printf("cannot encode a reply: %v", err)
|
|
return
|
|
}
|
|
timeout, cancel := context.WithTimeout(ctx, 10*time.Second)
|
|
defer cancel()
|
|
|
|
if err := s.channel.PublishWithContext(timeout, "", delivery.ReplyTo, false, false,
|
|
amqp.Publishing{
|
|
ContentType: "application/json",
|
|
CorrelationId: delivery.CorrelationId,
|
|
Body: body,
|
|
}); err != nil {
|
|
s.log.Printf("cannot reply to %s: %v", delivery.ReplyTo, err)
|
|
}
|
|
}
|