The mesh acts on what the catalogue decided a build meant
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.
This commit is contained in:
@@ -70,3 +70,23 @@ func eventID() (string, error) {
|
||||
// KeyModuleBuilt is what the builder announces when it has built something. The catalogue places
|
||||
// it in the module graph; nothing else need care.
|
||||
const KeyModuleBuilt = "module.builder.built"
|
||||
|
||||
// KeyModuleUpgraded is the catalogue saying a module's current version has moved.
|
||||
//
|
||||
// **The control plane hooks the meaning, not the build.** The builder says what it built; the
|
||||
// catalogue decides whether that was an upgrade — a rebuild producing the commit already current
|
||||
// is not one — and only this says anything the control plane can act on. Consuming the build
|
||||
// directly would make the control plane re-derive a decision another module already made, and the
|
||||
// two would eventually disagree (novox/hq ADR 0072).
|
||||
const KeyModuleUpgraded = "module.mesh-catalog.upgraded"
|
||||
|
||||
// UpgradeQueue is where those land. Durable and named, not a temporary queue: an upgrade announced
|
||||
// while the control plane is restarting is exactly the one that must not be missed.
|
||||
const UpgradeQueue = "control.upgrades"
|
||||
|
||||
// Upgraded is what the catalogue says when a module's current version moves.
|
||||
type Upgraded struct {
|
||||
Module string `json:"module"`
|
||||
Commit string `json:"commit"`
|
||||
Previous string `json:"previous"`
|
||||
}
|
||||
|
||||
@@ -50,6 +50,7 @@ type Server struct {
|
||||
listener Listener
|
||||
recorder Recorder
|
||||
log *log.Logger
|
||||
upgrader Upgrader
|
||||
}
|
||||
|
||||
// Records tells the server where to keep build results.
|
||||
@@ -59,6 +60,35 @@ type Server struct {
|
||||
// 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))
|
||||
@@ -88,6 +118,13 @@ func Connect(enroller Enroller, listener Listener) (*Server, error) {
|
||||
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)
|
||||
@@ -138,14 +175,39 @@ func (s *Server) Serve(ctx context.Context) error {
|
||||
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
|
||||
@@ -309,3 +371,35 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) {
|
||||
}
|
||||
_ = 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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user