diff --git a/cmd/mesh-control/main.go b/cmd/mesh-control/main.go index 2c8ef08..f69e90b 100644 --- a/cmd/mesh-control/main.go +++ b/cmd/mesh-control/main.go @@ -86,6 +86,8 @@ func run() error { return brokerCommand(args[1:]) case "serve": return serve(ctx) + case "upgrade": + return upgradeCommand(ctx, args[1:]) case "declare": return declare(ctx, args[1:]) case "overlay": @@ -140,6 +142,9 @@ func usage() { module forget remove one, unless a node runs it or the mesh holds things for it module forget --and-what-it-holds ...and discard its settings, secrets and ports too module issue --node a broker account for a module, scoped to its emits and consumes + upgrade what happens when this module's current version moves + upgrade roll-out [--together] ...send it to the machines running it + upgrade record ...record that they are behind, and send nothing status [--json] what is wrong, what is quiet, and what is out of date board [--listen ADDR] the same three questions, as a page that holds nothing api --issuer URL [--listen A] assign and unassign over http, for a surface that is not here diff --git a/cmd/mesh-control/push.go b/cmd/mesh-control/push.go index 0246def..e4f0c93 100644 --- a/cmd/mesh-control/push.go +++ b/cmd/mesh-control/push.go @@ -86,6 +86,12 @@ func serve(ctx context.Context) error { // And build results nobody was waiting for. A build triggered any other way than `build` // would otherwise be reported into the void, which is the same as not reporting it. server.Records(builds{inv}) + // And what the catalogue decided a build meant. The builder's own result is already handled + // above; this is the other half — the control plane is the only one of the three that knows + // which machines run the thing, so it is the one that acts (novox/hq ADR 0072). + if err := server.Follows(following{open}); err != nil { + return err + } return server.Serve(ctx) } diff --git a/cmd/mesh-control/upgrades.go b/cmd/mesh-control/upgrades.go new file mode 100644 index 0000000..f4ba58c --- /dev/null +++ b/cmd/mesh-control/upgrades.go @@ -0,0 +1,165 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "strings" + + "github.com/novox/mesh-control/internal/inventory" + "github.com/novox/mesh-control/internal/link" +) + +// following acts on what the catalogue announces. +// +// **It holds the stores, not a copy of the decision.** What to do about an upgrade is read when +// one arrives, so changing it takes effect on the next upgrade rather than on the next restart of +// the control plane. +type following struct{ open *stores } + +// Upgraded sends the machines running a module the version the catalogue now considers current — +// or records that they are behind, which is the default and needs no record. +// +// **Recording is not a second code path.** A machine that is not running what the mesh would send +// it is already something the mesh notices and reports; that is what `status` and `push --behind` +// are built on. So "record it" is the absence of an action, and the only thing this has to decide +// is whether to act. +func (f following) Upgraded(ctx context.Context, u link.Upgraded) error { + inv := f.open.inventory + + decision, err := inv.UpgradeOf(ctx, u.Module) + if err != nil { + return err + } + on, err := inv.Running(ctx, u.Module) + if err != nil { + return err + } + if len(on) == 0 { + fmt.Printf("%s moved to %s; no machine runs it\n", u.Module, shortCommit(u.Commit)) + return nil + } + if !decision.RollOut { + // Named rather than counted, and said even though nothing happens: an upgrade that was + // deliberately not rolled out and an upgrade that was never noticed look identical in a + // log that only speaks when it acts. + fmt.Printf("%s moved to %s; %s %s behind it, and this mesh records upgrades rather than "+ + "rolling them out — `push --behind` when you want them\n", + u.Module, shortCommit(u.Commit), readableList(on), isAre(len(on))) + return nil + } + + if decision.Together { + fmt.Printf("%s moved to %s; sending %s together\n", + u.Module, shortCommit(u.Commit), readableList(on)) + return sendTo(ctx, f.open, on) + } + // One at a time, and stopping at the first that fails. + // + // **Stopping is the point.** The machines are done one after another precisely so that a + // version that breaks the first one does not reach the rest; carrying on past a failure would + // make this the same as sending them together, only slower. + fmt.Printf("%s moved to %s; sending %s one at a time\n", + u.Module, shortCommit(u.Commit), readableList(on)) + for _, node := range on { + if err := sendTo(ctx, f.open, []string{node}); err != nil { + return fmt.Errorf("%s did not take %s, so the machines after it were left alone: %w", + node, u.Module, err) + } + } + return nil +} + +// readableList names machines the way a sentence does, because this is read by a person deciding +// whether an upgrade went where they expected. +func readableList(names []string) string { + switch len(names) { + case 0: + return "nothing" + case 1: + return names[0] + case 2: + return names[0] + " and " + names[1] + } + return strings.Join(names[:len(names)-1], ", ") + " and " + names[len(names)-1] +} + +func shortCommit(commit string) string { + if len(commit) > 8 { + return commit[:8] + } + return commit +} + +func isAre(n int) string { + if n == 1 { + return "is" + } + return "are" +} + +// upgradeCommand says what should happen when a module's current version moves. +func upgradeCommand(ctx context.Context, args []string) error { + set := flag.NewFlagSet("upgrade", flag.ContinueOnError) + together := set.Bool("together", false, + "send every machine running it at once, instead of one after another") + positionals, err := parseAround(set, args) + if err != nil { + return err + } + if len(positionals) == 0 { + return errors.New("upgrade [roll-out|record] [--together]") + } + module := positionals[0] + + open, err := openStores(ctx) + if err != nil { + return err + } + defer open.Close() + inv := open.inventory + + if len(positionals) == 1 { + decision, err := inv.UpgradeOf(ctx, module) + if err != nil { + return err + } + fmt.Println(sayUpgrade(module, decision)) + return nil + } + + var decision inventory.Upgrade + switch positionals[1] { + case "roll-out": + decision = inventory.Upgrade{RollOut: true, Together: *together} + case "record": + if *together { + // Refused rather than ignored: --together only means anything for a roll-out, and + // accepting it here would store a preference that never applies and looks like it does. + return errors.New("`--together` says how to roll out, so it cannot be given with " + + "`record`, which is the choice not to") + } + decision = inventory.Upgrade{} + default: + return fmt.Errorf("upgrade roll-out|record — not %q", positionals[1]) + } + if err := inv.SetUpgradeOf(ctx, module, decision); err != nil { + return err + } + fmt.Println(sayUpgrade(module, decision)) + return nil +} + +func sayUpgrade(module string, u inventory.Upgrade) string { + if !u.RollOut { + return fmt.Sprintf("when %s moves, the mesh records it and the machines running it are "+ + "reported as behind", module) + } + if u.Together { + return fmt.Sprintf("when %s moves, every machine running it is sent the new version "+ + "together", module) + } + return fmt.Sprintf("when %s moves, the machines running it are sent the new version one at "+ + "a time, stopping at the first that fails", module) +} diff --git a/internal/inventory/catalogue.go b/internal/inventory/catalogue.go index 90e4285..e29f994 100644 --- a/internal/inventory/catalogue.go +++ b/internal/inventory/catalogue.go @@ -796,3 +796,78 @@ func (i *Inventory) Catalogued(ctx context.Context) ([]Entry, error) { // providedBy is what the source column says for a module the control plane ships. const providedBy = "the control plane" + +// Upgrade is what the mesh decided to do when a module's current version moves. +type Upgrade struct { + // RollOut is true when the machines running it should be sent the new version. False means + // record it and stop — which needs no record of its own, because a machine not running what + // the mesh would send it is already something the mesh reports. + RollOut bool + // Together is true when every machine running it is sent the new version at once, rather than + // one after another. Only meaningful when RollOut is. + Together bool +} + +// UpgradeOf is what to do when this module moves. +// +// A module the mesh does not hold is not an error here: the catalogue may know of modules this +// mesh has never registered, and being told one of them moved is information, not a fault. The +// answer is the safe one — record it — because there is nothing to roll out to. +func (i *Inventory) UpgradeOf(ctx context.Context, module string) (Upgrade, error) { + var u Upgrade + var policy string + err := i.store.Pool().QueryRow(ctx, + `select upgrade, upgrade_together from module where name = $1`, module). + Scan(&policy, &u.Together) + if errors.Is(err, pgx.ErrNoRows) { + return Upgrade{}, nil + } + if err != nil { + return Upgrade{}, err + } + u.RollOut = policy == "roll-out" + return u, nil +} + +// SetUpgradeOf records what to do when this module moves. +func (i *Inventory) SetUpgradeOf(ctx context.Context, module string, u Upgrade) error { + policy := "record" + if u.RollOut { + policy = "roll-out" + } + tag, err := i.store.Pool().Exec(ctx, + `update module set upgrade = $2, upgrade_together = $3 where name = $1`, + module, policy, u.Together) + if err != nil { + return err + } + if tag.RowsAffected() == 0 { + return fmt.Errorf("this mesh holds no module called %s", module) + } + return nil +} + +// Running is every machine assigned a module, in a stable order. +// +// **Assigned, not reported.** A machine that is assigned the module and has not applied it yet is +// exactly the machine an upgrade most needs to reach; waiting for it to report the old version +// first would mean the machines furthest behind are the last to be caught up. +func (i *Inventory) Running(ctx context.Context, module string) ([]string, error) { + rows, err := i.store.Pool().Query(ctx, + `select n.name from assignment a join node n on n.id = a.node + where a.module = $1 order by n.name`, module) + if err != nil { + return nil, err + } + defer rows.Close() + + var out []string + for rows.Next() { + var name string + if err := rows.Scan(&name); err != nil { + return nil, err + } + out = append(out, name) + } + return out, rows.Err() +} diff --git a/internal/inventory/migrations/0021-what-to-do-when-a-module-is-upgraded.sql b/internal/inventory/migrations/0021-what-to-do-when-a-module-is-upgraded.sql new file mode 100644 index 0000000..d3a9acc --- /dev/null +++ b/internal/inventory/migrations/0021-what-to-do-when-a-module-is-upgraded.sql @@ -0,0 +1,26 @@ +-- What the mesh should do when the catalogue says a module has been upgraded. +-- +-- novox/hq ADR 0072. The builder announces a build, the catalogue decides whether that is an +-- upgrade, and the control plane is the only one of the three that knows which machines are +-- running the thing. So it is the one that acts -- and what it does has to be a decision somebody +-- made, not a behaviour compiled in. +-- +-- Two separate questions, deliberately not one: +-- +-- whether -- roll the new version out, or record that the machine is behind and stop there. +-- how -- one machine at a time, or all of them together. +-- +-- They are separate because the safe answer to the first is not the safe answer to the second: a +-- mesh may well want every upgrade applied automatically and still never want its only two +-- machines restarted in the same breath. +-- +-- Defaulted to recording rather than rolling out, and per module rather than mesh-wide. A mesh +-- that upgrades everything it builds the moment it builds it is a reasonable thing to want and a +-- terrible thing to arrive by default -- the first module to inherit it would be the control plane +-- itself, upgrading itself out from under the push that was applying it. +alter table module add column upgrade text not null default 'record' + check (upgrade in ('record', 'roll-out')); + +-- Only meaningful when upgrade is 'roll-out'. Kept anyway when it is not, so turning roll-out on +-- does not silently also decide this. +alter table module add column upgrade_together boolean not null default false; diff --git a/internal/link/events.go b/internal/link/events.go index d2266ff..7cb62fb 100644 --- a/internal/link/events.go +++ b/internal/link/events.go @@ -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"` +} diff --git a/internal/link/serve.go b/internal/link/serve.go index 74b0f8c..fef6712 100644 --- a/internal/link/serve.go +++ b/internal/link/serve.go @@ -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 +}