`declare` sends a node a signed declaration; `serve` now also consumes reports. Signed over the exact bytes published, which is what the node verifies. Anything re-encoding in between would sign one thing and check another, and a difference in key order alone would have a node refuse a declaration that was genuinely the mesh's. Sent to the node's queue directly rather than through the exchange: a declaration is for one node, and routing by name through a shared exchange means a binding per node that nothing removes when a node is retired. Enrolment now issues the node its own broker password, replacing the token's secret, and tells it the broker address, the fingerprint and the signing key -- so a node can reconnect after a restart without a person and a new token, which is what makes disconnection ordinary rather than a crisis. A report is a statement, not a write. What a node says it applied is its own account of its own machine, kept as a copy for recovery rather than as a source.
226 lines
7.7 KiB
Go
226 lines
7.7 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)
|
|
}
|
|
if err := channel.QueueBind(ControlQueue, KeyEnrol, Exchange, false, nil); err != nil {
|
|
conn.Close()
|
|
return nil, err
|
|
}
|
|
|
|
return &Server{conn: conn, channel: channel, enroller: enroller,
|
|
log: log.New(os.Stdout, "", log.LstdFlags)}, nil
|
|
}
|
|
|
|
func (s *Server) bindOrClose(channel *amqp.Channel, conn *amqp.Connection, key string) error {
|
|
if err := channel.QueueBind(ControlQueue, key, Exchange, false, nil); err != nil {
|
|
conn.Close()
|
|
return fmt.Errorf("cannot bind %s to %s/%s: %w", ControlQueue, Exchange, key, err)
|
|
}
|
|
return 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)
|
|
}
|
|
}
|