Files
mesh-controller/internal/link/serve.go
T
jschoubben 6c12780abe Describe the bus on its own terms
Comments framed the new bus by what it replaces — a comparison in almost
every explanation, which reads as though NATS were a variant of the old
thing rather than the mesh's nervous system. Removed throughout, and
OverAMQP becomes OverCurrent: the seam's two sides are the bus the mesh
runs on today and the one being built, not two protocols.

What remains is the client library's own package name, which is its name.
2026-09-26 23:51:00 +02:00

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, OverCurrent{Channel: 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
}