Three readers did not follow a moved foundation port (novox/hq 04-ISSUES/102),
and each took the control-node down in its own way: the control plane's own
store and broker connections, sealed at genesis with the port inside; and every
build the mesh ever recorded, kept as `<registry>:<port>/<module>/<artifact>@…`.
The control plane cannot open its own sealed connections to move a port, and it
cannot bind the store as a consumer would — a binding mints a credential. So its
settings get a third twin, `NAME_PORT`, read on top of the sealed value by the
store, the broker, the management API and the bus connection, and filled into
its container by a placeholder that names a seat, `${seat:mesh-store:5432}`,
from the node's given or mesh-assigned ports — never the manifest's number, and
empty when the mesh has nothing to add, so what genesis wrote stands. A value
that is still a placeholder is nothing said, aloud: the manifest naming it lands
in the next commit, once every control plane that composes it knows it.
A build is now recorded by digest and path — `artifact-store://<module>/<artifact>@…`
— and the store's address is composed in where a reference is used: the
declaration, the trust file, the bases a build is handed, a replay to the
catalogue. Over the network as `<node>.internal:<port>`; on the store's own node
before any network exists — every genesis push before its "network" step — by
loopback. A reference recorded before this, with an address, is re-routed the
same way when the mesh built it. The trust file and every provider's address
come from one derivation: the node's given port, over the mesh's assignment,
over the manifest's number.
novox/hq 04-ISSUES/102
699 lines
27 KiB
Go
699 lines
27 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"github.com/novox/mesh-controller/internal/envfile"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"log"
|
|
"os"
|
|
"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
|
|
replayer Replayer
|
|
// Messages the store could not take right now, held unacknowledged and tried again on a
|
|
// ticker, by subject (novox/hq issues 082, 083). again is the ticker's interval, zero meaning
|
|
// TryAgainAfter; giveUp is how long one is kept, zero meaning GiveUpAfter.
|
|
again time.Duration
|
|
giveUp time.Duration
|
|
parked map[string]*held
|
|
}
|
|
|
|
// held is one message the store could not take, kept to be tried again.
|
|
type held struct {
|
|
delivery amqp.Delivery
|
|
retry func(context.Context, amqp.Delivery)
|
|
what string
|
|
first time.Time
|
|
}
|
|
|
|
// ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and
|
|
// asking again, as when the store is restarting (novox/hq issue 082).
|
|
var ErrTryAgain = errors.New("not now, try again")
|
|
|
|
// TryAgainAfter is how often messages the store could not take are tried again. A store comes
|
|
// back in seconds, and a report a few seconds late is still current.
|
|
const TryAgainAfter = 2 * time.Second
|
|
|
|
// GiveUpAfter bounds how long one message is kept trying. A store that has not come back in this
|
|
// long is not restarting, and the message is let go with a line saying it was lost.
|
|
const GiveUpAfter = 2 * time.Minute
|
|
|
|
// Prefetch is how many messages the broker hands the control plane before it has settled them.
|
|
// More than one because a message the store could not take is held, unsettled, while the loop goes
|
|
// on answering others — an enrolment above all, which a host is waiting on (novox/hq issue 083).
|
|
// Bounded, because what is held is also what the broker has not kept on its own disk as pending.
|
|
const Prefetch = 64
|
|
|
|
// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to
|
|
// answer — an enrolment above all — while others wait for the store.
|
|
const PrefetchHeadroom = 8
|
|
|
|
// 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 — except the store unreachable for the moment, which is
|
|
// asked again for a bounded time (novox/hq issue 083).
|
|
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
|
|
}
|
|
|
|
// Answers binds the queue a catalogue's catch-up request arrives on.
|
|
//
|
|
// **Not bound unless something is listening**, for the same reason upgrades are not: a durable
|
|
// queue with no consumer fills quietly and the first symptom is a broker out of disk.
|
|
func (s *Server) Answers(r Replayer) error {
|
|
if _, err := s.channel.QueueDeclare(CatchUpQueue, true, false, false, false, nil); err != nil {
|
|
return fmt.Errorf("cannot declare the %s queue: %w", CatchUpQueue, err)
|
|
}
|
|
if err := s.channel.QueueBind(CatchUpQueue, KeyCatchingUp, EventsExchange, false, nil); err != nil {
|
|
return fmt.Errorf("cannot bind %s to %s/%s: %w", CatchUpQueue, EventsExchange, KeyCatchingUp, err)
|
|
}
|
|
s.replayer = r
|
|
return nil
|
|
}
|
|
|
|
// Connect opens the control plane's own connection to the broker.
|
|
//
|
|
// On the port MESH_BROKER_AMQP_PORT names when the node's settings moved the broker (novox/hq
|
|
// 04-ISSUES/102) — the URL is genesis's, sealed, and its port is the one thing in it the node may
|
|
// have moved since.
|
|
func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
|
url, err := envfile.Placed(AMQPVar)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
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 {
|
|
// A bounded prefetch rather than one. The loop still takes messages one at a time; what the
|
|
// prefetch buys is that a message the store could not take can be held while the loop goes on
|
|
// to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything held
|
|
// goes back to the broker if the control plane stops, because nothing held is acknowledged.
|
|
if err := s.channel.Qos(Prefetch, 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
|
|
}
|
|
}
|
|
|
|
// Its own queue and its own consumer, for the reason above: two consumers on one queue split
|
|
// its messages, and a catch-up request going to whichever half was not listening is a gap that
|
|
// looks like a working mesh.
|
|
var catchups <-chan amqp.Delivery
|
|
if s.replayer != nil {
|
|
catchups, err = s.channel.ConsumeWithContext(ctx, CatchUpQueue, "control-plane-catchup",
|
|
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)
|
|
}
|
|
|
|
again := s.again
|
|
if again == 0 {
|
|
again = TryAgainAfter
|
|
}
|
|
ticker := time.NewTicker(again)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil
|
|
case <-ticker.C:
|
|
if ctx.Err() != nil {
|
|
return nil
|
|
}
|
|
s.retryHeld(ctx)
|
|
case delivery, ok := <-catchups:
|
|
if !ok {
|
|
if catchups != nil {
|
|
return errors.New("the broker stopped delivering catch-up requests")
|
|
}
|
|
continue
|
|
}
|
|
s.catchingUp(ctx, delivery)
|
|
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(ctx, 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(ctx context.Context, 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
|
|
}
|
|
// A node's newer report supersedes one of its older reports still held: the older is its
|
|
// past, and recorded after the newer it would overwrite what the node is doing now.
|
|
subject := "report " + report.Node
|
|
s.supersede(subject, delivery)
|
|
if s.listener != nil {
|
|
if err := s.listener.Heard(context.Background(), report); err != nil {
|
|
// Held, not acknowledged, while the store cannot take it: the node reports an apply
|
|
// once, and a report lost here is a node the mesh never hears from again — the store
|
|
// restarting under the adoption that node just applied lost exactly that (issue 082).
|
|
what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared)
|
|
if s.tryLater(ctx, delivery, subject, what, err, s.handle) {
|
|
return
|
|
}
|
|
// 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))
|
|
}
|
|
s.settled(subject, delivery)
|
|
_ = delivery.Ack(false)
|
|
}
|
|
|
|
// tryLater holds a message the store could not take right now, to be tried again on the ticker,
|
|
// and says whether it did (novox/hq issues 082, 083).
|
|
//
|
|
// "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error the
|
|
// inventory reads as an outage. Anything else is an answer, and is left to the caller to settle.
|
|
// Held means unacknowledged and set aside: the loop goes on to the next message, so an enrolment a
|
|
// host is waiting on is answered while a report waits for the store. One message is held at most
|
|
// giveUpAfter; past it, it is let go with a line saying it was lost, and the caller settles it.
|
|
func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, subject, what string, err error,
|
|
retry func(context.Context, amqp.Delivery)) bool {
|
|
// Shutting down: nothing is settled. Unsettled, the broker hands the message to whatever
|
|
// consumes next — a cancelled context is not an answer about the message (issue 083, review).
|
|
if ctx.Err() != nil {
|
|
return true
|
|
}
|
|
if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) {
|
|
return false
|
|
}
|
|
if s.parked == nil {
|
|
s.parked = map[string]*held{}
|
|
}
|
|
h, ok := s.parked[subject]
|
|
if !ok || h.delivery.DeliveryTag != delivery.DeliveryTag {
|
|
// Held no further than the prefetch leaves room: past it, the broker would hand the loop
|
|
// nothing new — enrolments included — until something held was let go.
|
|
if !ok && len(s.parked) >= Prefetch-PrefetchHeadroom {
|
|
s.log.Printf("LOST %s: %d messages are already held for the store, and holding more "+
|
|
"would stop the queue: %v", what, len(s.parked), err)
|
|
return false
|
|
}
|
|
h = &held{delivery: delivery, retry: retry, what: what, first: time.Now()}
|
|
s.parked[subject] = h
|
|
s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err)
|
|
return true
|
|
}
|
|
if time.Since(h.first) >= s.giveUpAfter() {
|
|
delete(s.parked, subject)
|
|
s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err)
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
// supersede drops a message held for a subject when a newer one for it arrives: the older is
|
|
// acknowledged, because acting on it after the newer would undo the newer.
|
|
func (s *Server) supersede(subject string, newer amqp.Delivery) {
|
|
h, ok := s.parked[subject]
|
|
if !ok || h.delivery.DeliveryTag == newer.DeliveryTag {
|
|
return
|
|
}
|
|
delete(s.parked, subject)
|
|
s.log.Printf("set aside %s: a newer one arrived", h.what)
|
|
_ = h.delivery.Ack(false)
|
|
}
|
|
|
|
// settled forgets a message once it has been handled either way.
|
|
func (s *Server) settled(subject string, delivery amqp.Delivery) {
|
|
if h, ok := s.parked[subject]; ok && h.delivery.DeliveryTag == delivery.DeliveryTag {
|
|
delete(s.parked, subject)
|
|
}
|
|
}
|
|
|
|
// digest names a message by its content.
|
|
func digest(body []byte) string {
|
|
sum := sha256.Sum256(body)
|
|
return hex.EncodeToString(sum[:])
|
|
}
|
|
|
|
// retryHeld tries every held message again. Each handler holds it again, settles it, or lets it
|
|
// go past the bound.
|
|
func (s *Server) retryHeld(ctx context.Context) {
|
|
for _, h := range s.snapshot() {
|
|
if ctx.Err() != nil {
|
|
return
|
|
}
|
|
h.retry(ctx, h.delivery)
|
|
}
|
|
}
|
|
|
|
func (s *Server) snapshot() []*held {
|
|
out := make([]*held, 0, len(s.parked))
|
|
for _, h := range s.parked {
|
|
out = append(out, h)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (s *Server) giveUpAfter() time.Duration {
|
|
if s.giveUp == 0 {
|
|
return GiveUpAfter
|
|
}
|
|
return s.giveUp
|
|
}
|
|
|
|
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 {
|
|
request.Redelivered = delivery.Redelivered
|
|
accepted, err := s.enroller.Enrol(ctx, request)
|
|
switch {
|
|
case errors.Is(err, ErrTryAgain):
|
|
// Not a refusal: nothing was spent, and the same request asked again will be
|
|
// answered. Replied at once rather than held, so the node — which is waiting on
|
|
// this answer — decides when to ask, and the queue behind it moves (issue 083).
|
|
reply = EnrolReply{TryAgain: true, Refusal: "the mesh cannot answer right now; ask again"}
|
|
s.log.Printf("asked %q to enrol again shortly: %v", request.Node, err)
|
|
case err != nil && request.Redelivered:
|
|
// Said as what it most likely is: the broker handed this request over again after
|
|
// the control plane stopped mid-answer, and an enrolment already spent is not
|
|
// finished a second time. The node may need a new token.
|
|
s.log.Printf("refusing a redelivered enrolment for %q — it may have finished before "+
|
|
"the control plane stopped, and if the node did not get its answer it needs a new "+
|
|
"token: %v", request.Node, err)
|
|
case 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)
|
|
default:
|
|
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. Asked again by the same
|
|
// presenter, an enrolment finishes: the token is held for its key and spent last (issue 083).
|
|
_ = 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
|
|
}
|
|
// Each build result its own subject: none supersedes another, and recording one twice is
|
|
// harmless — the build is kept by its id.
|
|
subject := "build " + digest(delivery.Body)
|
|
s.supersede(subject, delivery)
|
|
if err := s.recorder.Built(ctx, result); err != nil {
|
|
// Held while the store cannot take it: a build result lost here is never announced (083).
|
|
if s.tryLater(ctx, delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) {
|
|
return
|
|
}
|
|
s.log.Printf("cannot keep a build result from %s: %v", result.On, err)
|
|
s.settled(subject, delivery)
|
|
_ = delivery.Reject(false)
|
|
return
|
|
}
|
|
s.settled(subject, delivery)
|
|
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)
|
|
}
|
|
|
|
// catchingUp answers a catalogue that has just started and may have missed builds.
|
|
//
|
|
// Acknowledged after the work. A replay that fails for a reason other than the store is not one
|
|
// that succeeds by being handed the same request again, so that is acknowledged and said; but a
|
|
// store that could not be read right now is held and asked again, bounded, rather than the
|
|
// request lost until the catalogue next restarts (issue 083). One request stands for all: a newer
|
|
// one supersedes one still held.
|
|
func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) {
|
|
const subject = "catch-up"
|
|
s.supersede(subject, delivery)
|
|
holding := false
|
|
defer func() {
|
|
if !holding {
|
|
s.settled(subject, delivery)
|
|
_ = delivery.Ack(false)
|
|
}
|
|
}()
|
|
if s.replayer == nil {
|
|
s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay")
|
|
return
|
|
}
|
|
announcements, err := s.replayer.Announceable(ctx)
|
|
if err != nil {
|
|
if s.tryLater(ctx, delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) {
|
|
holding = true
|
|
return
|
|
}
|
|
s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err)
|
|
return
|
|
}
|
|
sent := 0
|
|
for _, a := range announcements {
|
|
a.Replay = true
|
|
if err := EmitEvent(ctx, s.channel, KeyModuleBuilt, "control-plane", "", a); err != nil {
|
|
// Said and abandoned rather than retried: the catalogue asks again every time it
|
|
// starts, and half a graph delivered twice is no better than half delivered once.
|
|
s.log.Printf("replaying %s at %s failed, and the rest is abandoned: %v",
|
|
a.Module, short(a.Commit), err)
|
|
return
|
|
}
|
|
sent++
|
|
}
|
|
s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent)
|
|
}
|
|
|
|
// upgraded hands one announcement to whatever is following them.
|
|
//
|
|
// **Acknowledged whatever happens, but one thing.** A failure here is usually 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. The
|
|
// one exception is the upgrader saying the store could not be read for the moment (ErrTryAgain):
|
|
// that is held and asked again, bounded (novox/hq issue 083). Only the upgrader's word counts
|
|
// here, not an error that merely looks like an outage — a push that timed out on the second
|
|
// machine is not asked again, or the first would be pushed every few seconds for two minutes.
|
|
// A newer move of the same module supersedes one still held.
|
|
func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) {
|
|
var u Upgraded
|
|
_ = json.Unmarshal(delivery.Body, &u)
|
|
subject := "upgrade " + u.Module
|
|
s.supersede(subject, delivery)
|
|
holding := false
|
|
defer func() {
|
|
if !holding {
|
|
s.settled(subject, delivery)
|
|
_ = delivery.Ack(false)
|
|
}
|
|
}()
|
|
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 {
|
|
// Shutting down is not an answer about the announcement: left for the broker.
|
|
if ctx.Err() != nil {
|
|
holding = true
|
|
return
|
|
}
|
|
if errors.Is(err, ErrTryAgain) &&
|
|
s.tryLater(ctx, delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) {
|
|
holding = true
|
|
return
|
|
}
|
|
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
|
|
}
|