From 588aa424e20b0fdefd27e72331b9262164e14a1b Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 13 Sep 2026 01:55:07 +0200 Subject: [PATCH] The mesh acts on what the catalogue decided a build meant MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- cmd/mesh-control/main.go | 5 + cmd/mesh-control/push.go | 6 + cmd/mesh-control/upgrades.go | 165 ++++++++++++++++++ internal/inventory/catalogue.go | 75 ++++++++ ...1-what-to-do-when-a-module-is-upgraded.sql | 26 +++ internal/link/events.go | 20 +++ internal/link/serve.go | 94 ++++++++++ 7 files changed, 391 insertions(+) create mode 100644 cmd/mesh-control/upgrades.go create mode 100644 internal/inventory/migrations/0021-what-to-do-when-a-module-is-upgraded.sql 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 +}