The builder says what it built and the catalogue decides whether that was an upgrade. Only the control plane knows which machines run the thing, so it is the one that acts — and what it does is a choice somebody recorded, not a behaviour compiled in: record that they are behind, or send it, one machine at a time or together. Recording is the absence of an action rather than a second path: a machine not running what the mesh would send it is already something the mesh reports. Defaulted to recording. A mesh that rolls out everything it builds the moment it builds it is reasonable to want and a bad thing to arrive by default — the first module to inherit it would be the control plane, upgrading itself out from under the push applying it.
406 lines
15 KiB
Go
406 lines
15 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
|
|
}
|
|
|
|
// Recorder keeps what builders say.
|
|
//
|
|
// Separate from Listener because they are different things arriving: a report is a node saying
|
|
// what it did with a declaration, and a build result is a machine saying what came of some work.
|
|
// One interface carrying both would mean an implementation of one having to say something about
|
|
// the other.
|
|
type Recorder interface {
|
|
Built(ctx context.Context, result BuildResult) error
|
|
}
|
|
|
|
type Server struct {
|
|
conn *amqp.Connection
|
|
channel *amqp.Channel
|
|
enroller Enroller
|
|
listener Listener
|
|
recorder Recorder
|
|
log *log.Logger
|
|
upgrader Upgrader
|
|
}
|
|
|
|
// Records tells the server where to keep build results.
|
|
//
|
|
// Set after Connect rather than passed to it, because a control plane that only publishes — the
|
|
// `build` command, which waits for its own answer — needs a connection and no recorder, and
|
|
// making it supply one would have it construct something it never uses.
|
|
func (s *Server) Records(r Recorder) { s.recorder = r }
|
|
|
|
// Upgrader is what the control plane does when the catalogue says a module moved.
|
|
//
|
|
// An interface for the same reason Enroller is one: deciding what an upgrade means for the
|
|
// machines running it is a different concern from noticing that one was announced, and only the
|
|
// first needs a database.
|
|
type Upgrader interface {
|
|
// Upgraded is told which module moved and between which commits. An error is logged and the
|
|
// message is not requeued: an upgrade the control plane could not act on is not one it will
|
|
// act on by being handed the same message again, and a poison message on a durable queue
|
|
// would stop every upgrade behind it.
|
|
Upgraded(ctx context.Context, u Upgraded) error
|
|
}
|
|
|
|
// Follows says what to do about upgrades, and binds the queue they arrive on.
|
|
//
|
|
// **Not bound unless something is listening.** A durable queue bound to every upgrade with no
|
|
// consumer fills up quietly, and the first symptom is a broker out of disk rather than anything
|
|
// about modules.
|
|
func (s *Server) Follows(u Upgrader) error {
|
|
if _, err := s.channel.QueueDeclare(UpgradeQueue, true, false, false, false, nil); err != nil {
|
|
return fmt.Errorf("cannot declare the %s queue: %w", UpgradeQueue, err)
|
|
}
|
|
if err := s.channel.QueueBind(UpgradeQueue, KeyModuleUpgraded, EventsExchange, false, nil); err != nil {
|
|
return fmt.Errorf("cannot bind %s to %s/%s: %w", UpgradeQueue, EventsExchange, KeyModuleUpgraded, err)
|
|
}
|
|
s.upgrader = u
|
|
return nil
|
|
}
|
|
|
|
// 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)
|
|
}
|
|
// The events exchange too. The control plane is not the only publisher on it — modules
|
|
// announce onto it with their own accounts — but it is the only thing permitted to create it,
|
|
// for the same reason it is the only thing permitted to create the direct one.
|
|
if err := channel.ExchangeDeclare(EventsExchange, "topic", true, false, false, false, nil); err != nil {
|
|
conn.Close()
|
|
return nil, fmt.Errorf("cannot declare the %s exchange: %w", EventsExchange, 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, KeyAlive, KeyBuilt} {
|
|
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
|
|
}
|
|
|
|
// The upgrade queue, when something is listening for them. A second queue rather than a
|
|
// second consumer on the first: two consumers on one queue split its messages between them,
|
|
// which is the fault the comment above exists about. Two queues share nothing.
|
|
var upgrades <-chan amqp.Delivery
|
|
if s.upgrader != nil {
|
|
upgrades, err = s.channel.ConsumeWithContext(ctx, UpgradeQueue, "control-plane-upgrades",
|
|
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,%s,%s}",
|
|
ControlQueue, Exchange, KeyEnrol, KeyReport, KeyAlive, KeyBuilt)
|
|
if s.upgrader != nil {
|
|
s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded)
|
|
}
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case delivery, ok := <-upgrades:
|
|
// A nil channel blocks for ever, so this case simply never fires when nothing is
|
|
// listening for upgrades. Closed is different, and means the broker stopped.
|
|
if !ok {
|
|
if upgrades != nil {
|
|
return errors.New("the broker stopped delivering upgrades")
|
|
}
|
|
continue
|
|
}
|
|
s.upgraded(ctx, delivery)
|
|
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)
|
|
case KeyAlive:
|
|
s.handleAlive(delivery)
|
|
case KeyBuilt:
|
|
s.handleBuilt(ctx, 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)
|
|
}
|
|
}
|
|
|
|
// handleAlive records that a node was heard from, and nothing else.
|
|
//
|
|
// Deliberately silent: a node saying it is there every minute would fill the log with the
|
|
// ordinary case, and a log where the ordinary case is loud is a log nobody reads.
|
|
func (s *Server) handleAlive(delivery amqp.Delivery) {
|
|
var alive Alive
|
|
if err := json.Unmarshal(delivery.Body, &alive); err != nil || alive.Node == "" {
|
|
_ = delivery.Reject(false)
|
|
return
|
|
}
|
|
if s.listener != nil {
|
|
if err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil {
|
|
s.log.Printf("could not record that %s is here: %v", alive.Node, err)
|
|
}
|
|
}
|
|
_ = delivery.Ack(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)
|
|
}
|
|
}
|
|
|
|
// handleBuilt 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
|
|
// to the asker; one triggered any other way is published here, and without this it would be
|
|
// reported into the void — which is the same as not reporting it.
|
|
func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) {
|
|
var result BuildResult
|
|
if err := json.Unmarshal(delivery.Body, &result); err != nil {
|
|
s.log.Printf("a build result could not be read: %v", err)
|
|
_ = delivery.Reject(false)
|
|
return
|
|
}
|
|
if s.recorder == nil {
|
|
// Nothing to keep it in. Rejected rather than dropped silently, so the broker's own
|
|
// counters show something arriving that nothing handles.
|
|
s.log.Printf("a build result arrived and this control plane keeps none")
|
|
_ = delivery.Reject(false)
|
|
return
|
|
}
|
|
if err := s.recorder.Built(ctx, result); err != nil {
|
|
s.log.Printf("cannot keep a build result from %s: %v", result.On, err)
|
|
_ = delivery.Reject(false)
|
|
return
|
|
}
|
|
switch {
|
|
case result.Failed != "":
|
|
s.log.Printf("%s could not build %s", result.On, result.Repository)
|
|
default:
|
|
s.log.Printf("%s built %s from %s", result.On, result.Repository, result.Commit)
|
|
}
|
|
_ = delivery.Ack(false)
|
|
}
|
|
|
|
// upgraded hands one announcement to whatever is following them.
|
|
//
|
|
// **Acknowledged whatever happens.** A failure here is the control plane being unable to act on an
|
|
// upgrade — a machine that cannot be resolved, a broker that will not take a declaration — and
|
|
// none of those get better by being handed the same message again. Requeuing would put a poison
|
|
// message at the head of a durable queue and stop every upgrade behind it, which turns one module
|
|
// nobody can push into a mesh that stops following its own catalogue.
|
|
func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) {
|
|
defer func() { _ = delivery.Ack(false) }()
|
|
var u Upgraded
|
|
if err := json.Unmarshal(delivery.Body, &u); err != nil {
|
|
s.log.Printf("an upgrade announcement could not be read: %v", err)
|
|
return
|
|
}
|
|
if u.Module == "" {
|
|
s.log.Printf("an upgrade announcement named no module; ignored")
|
|
return
|
|
}
|
|
if err := s.upgrader.Upgraded(ctx, u); err != nil {
|
|
s.log.Printf("%s moved to %s and the mesh could not act on it: %v",
|
|
u.Module, short(u.Commit), err)
|
|
}
|
|
}
|
|
|
|
// short is a commit as people read it.
|
|
func short(commit string) string {
|
|
if len(commit) > 8 {
|
|
return commit[:8]
|
|
}
|
|
return commit
|
|
}
|