Files
mesh-controller/internal/link/serve.go
T
jschoubben 8b974deb42 A working private network, and four reasons it did not work
Three machines across two sites, two of them behind no reachable address, all
nine paths open. The mesh computes the graph, delivers it as a declaration, and
the nodes bring it up.

Every fault below looked like success from inside the mesh: the graph was
right, the files were right, the services were up, every node reported it had
applied. None was reachable by reasoning.

A running interface does not re-read its configuration. A node joins, every
existing node's peer list changes, the file is replaced -- and the service is
already running, so nothing reloads it. Fixed as declared state rather than a
command: the service must reflect the file. A command to restart would be an
action, and the link may not carry one. The host refused exactly that, which is
how this shape was arrived at.

A hub sharing a site with a spoke appeared twice in that spoke's peer list --
once as a direct peer, once as the route of last resort. WireGuard takes one
entry per key and refuses the file. The ordinary shape of a small mesh, and in
none of the tests written before it ran.

Two nodes at one site that neither can be dialled were peered directly. Nobody
opens the path, and the direct route is more specific than the hub's, so it
wins and blackholes -- this design's own warning arriving in its
implementation. They now route through the hub unless one end can be dialled.

And Docker sets the FORWARD policy to DROP, so a hub with ip_forward enabled
carried nothing between its spokes. The substrate at tier 1 silently breaks the
network at tier 2, and nothing in either tier's state says so. The hub inserts
its own rule above those chains and removes it on the way down.

Two weak tests found by injection along the way: one asserted the keepalive
rule only against the hub, whose peer entries happen not to set that field at
all, so it tested an absence; the other checked the firewall rules by looking
for FORWARD anywhere, which the PostDown line satisfies on its own.
2026-08-29 18:04:15 +02:00

238 lines
8.4 KiB
Go

package link
import (
"context"
"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, request EnrolRequest) (EnrolReply, error)
}
// Server consumes what nodes say.
// Listener is what the control plane does with a report. Separate from Enroller so the two can
// be given independently, and so a server that only sends declarations needs neither.
type Listener interface {
Heard(ctx context.Context, report Report) error
}
type Server struct {
conn *amqp.Connection
channel *amqp.Channel
enroller Enroller
listener Listener
log *log.Logger
}
// Connect opens the control plane's own connection to the broker.
func Connect(enroller Enroller, listener Listener) (*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, listener: listener,
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
}
if s.listener != nil {
if err := s.listener.Heard(context.Background(), report); err != nil {
// Said rather than swallowed. A report the mesh heard and failed to write down is a
// node whose recovery copy is silently older than it looks.
s.log.Printf("could not record %s's report: %v", report.Node, err)
}
}
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)
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)
}
}