diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index 2724459..5f5fad5 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -56,6 +56,14 @@ is dialled except the broker. MESH_REGISTRY host:port to publish artifacts to, when the mesh has not said MESH_BINDING a file the mesh wrote saying where the artifact store is MESH_WORKSPACE where to clone and build (default: a temporary directory) + +It also builds one module and stops, which is how a mesh is raised — before there is a +broker to take work from or a registry to publish into: + + mesh-builder build [--path P] [--ref COMMIT] [--registry HOST:PORT] + +Without --registry the artifacts stay in this machine's container runtime, named by the +digest of their own configuration. The result is printed as JSON. ` func run() error { @@ -64,6 +72,8 @@ func run() error { case "version": fmt.Println(version) return nil + case "build": + return buildOnce(context.Background(), os.Args[2:]) default: fmt.Print(usage) return nil @@ -82,9 +92,18 @@ func run() error { if workspace == "" { workspace = os.TempDir() + "/mesh-builder" } - on, err := os.Hostname() - if err != nil { - on = "a build machine" + // **The mesh's name for this machine, not the container's.** A build is reported to the rest + // of the mesh, and a report whose origin reads `104cb10e105b` names something no other module + // can look up. The mesh already knows the answer and has a way to say it — `${machine:name}` + // in the environment file this module is handed — so the hostname is only what is left when + // nobody said. + on := os.Getenv("MESH_NODE") + if on == "" { + hostname, err := os.Hostname() + if err != nil { + return fmt.Errorf("this build machine has no name: nothing said MESH_NODE and the host would not say either: %w", err) + } + on = hostname } ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) @@ -183,6 +202,7 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis Name: made.Name, Kind: made.Kind, Reference: made.Reference, }) } + result.Against = built.Against fmt.Printf(" built %s from %s\n", built.Manifest.Module, short(built.Commit)) } } @@ -211,11 +231,45 @@ func answer(ctx context.Context, channel *amqp.Channel, publisher builder.Publis }); err != nil { fmt.Fprintf(os.Stderr, "cannot answer a build request: %v\n", err) } + // **And announced, which is a different act from answering.** The reply goes to whoever asked + // and is correlated to their request; this says to the whole mesh that a module now exists at + // a commit, and the catalogue places it in the module graph (novox/hq ADR 0072). A build + // nobody asked for still has to be announced, or the graph knows less than the registry does. + // + // Only on success: a failed build produced no module-version, and announcing one would put + // something in the graph that was never made. + if result.Failed == "" && result.Commit != "" { + announced := map[string]any{ + "module": moduleOf(result.Manifest), "commit": result.Commit, + "repository": result.Repository, "path": result.Path, "ref": result.Ref, + "manifest": json.RawMessage(result.Manifest), "against": result.Against, + "made": result.Made, + } + if err := link.EmitEvent(publishCtx, channel, link.KeyModuleBuilt, "builder", on, announced); err != nil { + // Said, not fatal: the build happened and was answered. A module the catalogue has not + // heard of is a gap somebody can close; a build reported as failed because announcing + // it failed is a lie about work that was done. + fmt.Fprintf(os.Stderr, " built, but could not announce it: %v\n", err) + } + } + // Acknowledged only once the answer is away, so a builder that dies before answering leaves // the request for another machine rather than losing it. _ = delivery.Ack(false) } +// moduleOf reads the module's name out of the manifest it just built, which is the only place it is +// authoritative — the request named a repository and a path, not a module. +func moduleOf(manifest json.RawMessage) string { + var named struct { + Module string `json:"module"` + } + if err := json.Unmarshal(manifest, &named); err != nil { + return "" + } + return named.Module +} + func short(commit string) string { if len(commit) > 8 { return commit[:8] diff --git a/cmd/mesh-builder/once.go b/cmd/mesh-builder/once.go new file mode 100644 index 0000000..8a98042 --- /dev/null +++ b/cmd/mesh-builder/once.go @@ -0,0 +1,134 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "flag" + "fmt" + "os" + + "github.com/novox/mesh-control/internal/builder" +) + +// buildOnce is the builder doing one build and stopping, with no broker and no mesh. +// +// **This is how a mesh is raised** (novox/hq ADR 0073). The installer carries this program and runs +// it once, before anything else exists, to produce the control plane from the same repository and +// path that every later rebuild of it will use. What raises the mesh is therefore the same thing +// that maintains it — not a second mechanism that has to be kept in step with the first and is +// exercised once per new mesh, which is how often enough to rot. +// +// It takes no work from a queue and answers nobody: there is no broker yet, and the only thing +// waiting for the answer is the installer that started it. So the result goes to standard output as +// JSON, which is what the installer reads. +// +// With no --registry it publishes nowhere and the image stays in this machine's container runtime, +// named by the digest of its own configuration (see builder.Local). That is the genesis case. With +// one, it publishes as it always does — the same command is how an operator builds a module by hand +// on a mesh that already exists. +func buildOnce(ctx context.Context, args []string) error { + set := flag.NewFlagSet("build", flag.ContinueOnError) + path := set.String("path", "", "the module's directory inside the repository (default: its root)") + ref := set.String("ref", "", "the commit to build; a branch is a moving target somebody else controls") + registry := set.String("registry", "", + "host:port to publish to. Without it the artifacts stay in this machine's container runtime, which is the genesis case") + workspace := set.String("workspace", "", "where to clone and build (default: a temporary directory)") + positionals, err := parseAround(set, args) + if err != nil { + return err + } + if len(positionals) != 1 { + return errors.New("mesh-builder build [--path P] [--ref COMMIT] [--registry HOST:PORT]") + } + repository := positionals[0] + + where := *workspace + if where == "" { + where = os.TempDir() + "/mesh-builder-once" + } + + // Local unless told otherwise, because the moment this exists for has nowhere to publish. A + // default pointing at a registry would mean genesis failing at a push to something that is not + // there yet, one step away from the thing that could explain it. + var publisher builder.Publisher = builder.Local{Run: builder.Command} + if *registry != "" { + publisher = builder.Registry{Address: *registry, Run: builder.Command} + } + + fmt.Fprintf(os.Stderr, "building %s", repository) + if *path != "" { + fmt.Fprintf(os.Stderr, " at %s", *path) + } + if *ref != "" { + fmt.Fprintf(os.Stderr, " at %s", *ref) + } + fmt.Fprintln(os.Stderr) + + built, buildErr := builder.Build(ctx, builder.Command, publisher, repository, *path, *ref, where) + if buildErr != nil { + return buildErr + } + + // To standard output, and everything else to standard error, so the caller can read this + // without having to separate it from progress. + out := onceResult{ + Module: built.Manifest.Module, + Commit: built.Commit, + Repository: repository, + Path: *path, + Ref: *ref, + Manifest: built.Manifest, + Against: built.Against, + } + for _, made := range built.Built { + out.Made = append(out.Made, madeArtifact{Name: made.Name, Kind: made.Kind, Reference: made.Reference}) + } + body, err := json.MarshalIndent(out, "", " ") + if err != nil { + return err + } + fmt.Println(string(body)) + fmt.Fprintf(os.Stderr, " built %s from %s\n", built.Manifest.Module, short(built.Commit)) + return nil +} + +// onceResult is what one build reports to whoever started it. +// +// The same fields the mesh records for a build, so a reader comparing a genesis build against an +// ordinary one is comparing the same thing said the same way. +type onceResult struct { + Module string `json:"module"` + Commit string `json:"commit"` + Repository string `json:"repository"` + Path string `json:"path,omitempty"` + Ref string `json:"ref,omitempty"` + Manifest any `json:"manifest"` + Made []madeArtifact `json:"made"` + Against []string `json:"against,omitempty"` +} + +type madeArtifact struct { + Name string `json:"name"` + Kind string `json:"kind"` + Reference string `json:"reference"` +} + +// parseAround lets flags appear on either side of the repository, because a person writing this by +// hand will put them wherever reads best and the standard parser stops at the first thing that is +// not a flag. The same helper the control plane's commands use, for the same reason. +func parseAround(set *flag.FlagSet, args []string) ([]string, error) { + var positionals []string + rest := args + for { + if err := set.Parse(rest); err != nil { + return nil, err + } + rest = set.Args() + if len(rest) == 0 { + return positionals, nil + } + positionals = append(positionals, rest[0]) + rest = rest[1:] + } +} 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/broker/management.go b/internal/broker/management.go index 5dbc3ae..0277d61 100644 --- a/internal/broker/management.go +++ b/internal/broker/management.go @@ -242,11 +242,14 @@ func (m *Management) CreateBuilderAccount(ctx context.Context, name, password st // and a queue nobody may declare is a queue that exists only if the control plane has // already run, which makes the order they start in matter. "configure": "^" + builds + "$", - // The exchange, and nothing else. **Not the default exchange**: permission there is + // Two exchanges, and nothing else. **Not the default exchange**: permission there is // granted per exchange rather than per queue, so a builder allowed to use it could // publish into any node's queue — the privilege a build machine most obviously should - // not have. Answers go through the exchange, which is why they can. - "write": "^" + regexp.QuoteMeta(ExchangeName) + "$", + // not have. Answers go through the node exchange; announcing what was built goes through + // the events exchange, which is a different act with a different audience (novox/hq + // ADR 0072). A builder that could answer and not announce would leave the module graph + // knowing less than the registry does. + "write": "^(" + regexp.QuoteMeta(ExchangeName) + "|" + regexp.QuoteMeta(EventsExchangeName) + ")$", // The build queue and nothing else. Not another machine's declarations. "read": "^" + builds + "$", }); err != nil { diff --git a/internal/builder/builder.go b/internal/builder/builder.go index bc1ba2b..ee5e9f4 100644 --- a/internal/builder/builder.go +++ b/internal/builder/builder.go @@ -11,6 +11,7 @@ import ( "os" "os/exec" "path/filepath" + "regexp" "sort" "strings" @@ -42,6 +43,15 @@ type Publisher interface { // Result is everything one build produced. type Result struct { + // Against is every pinned image this build was built on top of, read out of its own inputs. + // + // **Derived, not declared** (novox/hq ADR 0009): a declared list of dependencies drifts from + // what the code actually uses, and an artifact is out of date when anything it was built + // against moved. These are artifact references rather than module-versions, because that is + // what a build input names; resolving them to modules is the catalogue's work, since it is + // what knows which module-version published which artifact. + Against []string + // Manifest is the module as the mesh should hold it: artifacts resolved to digests. Manifest catalogue.Manifest // Commit is what was built, so "is this current?" is answerable without building again. @@ -124,7 +134,8 @@ func Build(ctx context.Context, run Runner, publish Publisher, if err != nil { return Result{}, err } - return Result{Manifest: resolved, Commit: commit, Built: built}, nil + return Result{Manifest: resolved, Commit: commit, Built: built, + Against: against(within, manifest)}, nil } // inside resolves a module's path within a clone, and refuses one that leaves it. @@ -157,6 +168,39 @@ func describe(path string) string { return path } +// pinnedImage matches an image reference pinned by digest, which is the only kind a build input is +// allowed to name — a tag is something somebody else can move under you. +var pinnedImage = regexp.MustCompile(`[A-Za-z0-9][A-Za-z0-9._/:-]*@sha256:[0-9a-f]{64}`) + +// against reads what this module's image artifacts are built on top of, out of the files that +// build them. Nothing is guessed: a reference that is not written down is not reported. +func against(within string, manifest catalogue.Manifest) []string { + if manifest.Build == nil { + return nil + } + seen := map[string]bool{} + var out []string + for _, a := range manifest.Build.Artifacts { + if a.Kind != catalogue.ArtifactImage || a.From == "" { + continue + } + body, err := os.ReadFile(filepath.Join(within, a.From)) + if err != nil { + // Not fatal: the build itself already failed if this file was needed and missing, and + // reporting no edges is honest where inventing them would not be. + continue + } + for _, found := range pinnedImage.FindAllString(string(body), -1) { + if !seen[found] { + seen[found] = true + out = append(out, found) + } + } + } + sort.Strings(out) + return out +} + // ManifestName is the one file a module repository must have. // // At the root, and named the same in every repository. A convention somebody can look for beats a diff --git a/internal/builder/local.go b/internal/builder/local.go new file mode 100644 index 0000000..7e471db --- /dev/null +++ b/internal/builder/local.go @@ -0,0 +1,61 @@ +package builder + +import ( + "context" + "fmt" + "strings" +) + +// Local is what a build publishes into when there is nowhere to publish yet. +// +// **This exists for exactly one moment: raising a mesh** (novox/hq ADR 0073). The installer builds +// the control plane on the machine that is about to run it, and at that moment there is no registry +// — the registry is installed afterwards, by the control plane this build produces. So the artifact +// stays where the build left it: in the machine's own container runtime. +// +// That is not a weaker kind of pinning. An image held locally is named by the digest of its own +// configuration, which is content-addressed and unforgeable and requires nothing to have served it +// — the same identity the installer has always used for the image it carried. What changes when a +// registry exists is not that the artifact becomes exact, but that something other than this +// machine can fetch it. +// +// It is deliberately unable to publish an archive. An archive has no local identity to fall back +// on: it is bytes that only mean something once something serves them at a URL. A build that +// produces one before there is anywhere to put it has produced nothing usable, and saying so is +// better than returning a path on a disk that no other machine can read. +type Local struct { + // Run is how docker is invoked, so a test does not need one. + Run Runner +} + +// PublishImage leaves the image where the build put it, and names it by its own configuration. +// +// The local tag is not returned: a tag is a name somebody can move, and every other reference in a +// resolved manifest is exact. The digest is read back from the runtime rather than computed, for +// the same reason the registry publisher reads it back from the registry — what matters is what +// will be served for that reference, and only the thing serving it can say. +func (l Local) PublishImage(ctx context.Context, localTag, repository string) (string, error) { + out, err := l.Run(ctx, "", "docker", "image", "inspect", "--format", "{{.Id}}", localTag) + if err != nil { + return "", fmt.Errorf("cannot read back the image just built as %s: %w", localTag, err) + } + id := strings.TrimSpace(out) + if !strings.HasPrefix(id, "sha256:") || len(id) != len("sha256:")+64 { + // Refused rather than passed on. A bundle naming an image by anything a person could move + // is refused by the machine applying it, and a value that is not an identity would fail + // there instead — one step further from the thing that could explain it. + return "", fmt.Errorf( + "the container runtime named the image just built %q, which is not an image id: "+ + "sha256 and sixty-four hex characters", id) + } + return id, nil +} + +// PublishArchive refuses, and says why rather than inventing somewhere to put it. +func (l Local) PublishArchive(ctx context.Context, repository string, body []byte, digest string) (string, error) { + return "", fmt.Errorf( + "%s declares an archive, and this build has nowhere to publish one. An image can stay in "+ + "the machine's own runtime and still be named exactly; an archive is bytes that mean "+ + "nothing until something serves them. Build this once the mesh has a registry", + repository) +} diff --git a/internal/catalogue/declaration.go b/internal/catalogue/declaration.go index a8f30bf..3ead392 100644 --- a/internal/catalogue/declaration.go +++ b/internal/catalogue/declaration.go @@ -456,6 +456,11 @@ func (r Resolution) Declaration(with Rendering) ([]map[string]any, error) { if err := pinned(copied, m.Module); err != nil { return nil, err } + // And the other half of the same question: an image nobody has published yet is named + // by an artifact rather than by a placeholder digest, and is just as unrunnable. + if err := built(copied, m.Module); err != nil { + return nil, err + } publishedOn(copied, m.Module, with) copied["id"] = m.Module + "." + fmt.Sprint(resource["id"]) // A service saying what it reflects names resources within its own module, so those @@ -887,6 +892,30 @@ func pinned(resource map[string]any, module string) error { module, resource["id"]) } +// built refuses a container that still names an artifact nobody has made. +// +// **Said here, by the thing that knows what "artifact" means.** A container naming an artifact is +// a module saying "the mesh builds this"; the field is resolved into an image when a build +// publishes one, and until then there is nothing to run. A host receiving it refuses the whole +// declaration — correctly, since its language has no such field — but what it can say is that a +// container does not use "artifact", which tells a reader the manifest is malformed. It is not: +// it is unbuilt, which is a different problem with a different fix, and only the mesh is in a +// position to tell them apart (novox/hq ADR 0073). +func built(resource map[string]any, module string) error { + if fmt.Sprint(resource["type"]) != "container" { + return nil + } + artifact, ok := resource["artifact"].(string) + if !ok || artifact == "" { + return nil + } + return fmt.Errorf( + "%s has not been built for this mesh: %v is its %q artifact, and no build has published "+ + "one. Build it — `build --path ` — and the module will name what "+ + "came out instead", + module, resource["id"], artifact) +} + // publishedOn puts the machine's own port on the outside of a container's mapping. // // **A module writes the port the software uses; the mesh says where the machine puts it** 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/build.go b/internal/link/build.go index 2b44228..64ecd96 100644 --- a/internal/link/build.go +++ b/internal/link/build.go @@ -82,6 +82,10 @@ type BuildResult struct { // Made is each artifact, for reporting. Made []MadeArtifact `json:"made,omitempty"` + // Against is every pinned image this was built on top of, read out of the build's own inputs + // (novox/hq ADR 0009). The catalogue turns these into edges; nothing else need care. + Against []string `json:"against,omitempty"` + // Failed is why, when it did. Failed string `json:"failed,omitempty"` } diff --git a/internal/link/events.go b/internal/link/events.go new file mode 100644 index 0000000..7cb62fb --- /dev/null +++ b/internal/link/events.go @@ -0,0 +1,92 @@ +package link + +import ( + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "fmt" + "time" + + amqp "github.com/rabbitmq/amqp091-go" +) + +// Emitting a module event from Go. +// +// **Every event rides one topic exchange** (novox/hq ADR 0042), which is not the direct exchange +// nodes and the control plane speak over. A module that announces something publishes here, and +// consumers bind their own durable queue to a pattern over it. +// +// This exists because the builder is a module written in Go while every other emitter is +// TypeScript on the sdk. The envelope is the sdk's, reproduced exactly: the body is the payload +// alone and everything about the event travels as headers. A second shape would be a second thing +// for consumers to handle, and they are written against the first. +const ( + // EventsExchange is where every event rides. Named here rather than imported from the broker + // package for the same reason BuildQueueName is duplicated there — one direction of dependency. + EventsExchange = "mesh.events" +) + +// EmitEvent publishes one module event, in the envelope the sdk's consumers expect. +// +// Persistent, because an event that a broker restart loses is not an announcement. The publish is +// not confirmed here: the caller has already done the work the event describes, and a build that +// succeeded must not be reported as failed because saying so failed. +func EmitEvent(ctx context.Context, channel *amqp.Channel, eventType, source, node string, body any) error { + payload, err := json.Marshal(body) + if err != nil { + return fmt.Errorf("cannot serialise a %s event: %w", eventType, err) + } + id, err := eventID() + if err != nil { + return err + } + return channel.PublishWithContext(ctx, EventsExchange, eventType, false, false, amqp.Publishing{ + ContentType: "application/json", + DeliveryMode: amqp.Persistent, + MessageId: id, + Timestamp: time.Now().UTC(), + Body: payload, + Headers: amqp.Table{ + "x-event-id": id, + "x-source": source, + "x-node": node, + "x-time": time.Now().UTC().Format(time.RFC3339), + "content-type": "application/json", + }, + }) +} + +// eventID is what a consumer deduplicates on: delivery is at-least-once, so a handler must be able +// to tell a redelivery from a second event, and only the emitter can say which it is. +func eventID() (string, error) { + raw := make([]byte, 16) + if _, err := rand.Read(raw); err != nil { + return "", fmt.Errorf("cannot make an event id: %w", err) + } + return hex.EncodeToString(raw), nil +} + +// 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 +} diff --git a/module.json b/module.json new file mode 100644 index 0000000..a6b7c36 --- /dev/null +++ b/module.json @@ -0,0 +1,68 @@ +{ + "module": "mesh-control", + "version": "1", + "slug": "control", + "capabilities": [ + "container-runtime" + ], + "claims": [ + { + "name": "the-control-plane", + "scope": "mesh" + } + ], + "own-secrets": { + "inventory": "/var/lib/mesh/mesh-control/inventory", + "identity": "/var/lib/mesh/mesh-control/identity", + "licences": "/var/lib/mesh/mesh-control/licences", + "broker": "/var/lib/mesh/mesh-control/broker", + "broker-management": "/var/lib/mesh/mesh-control/broker-management", + "broker-address": "/var/lib/mesh/mesh-control/broker-address" + }, + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/mesh-control", + "mode": "0700" + }, + { + "id": "control-env", + "type": "file", + "path": "/var/lib/mesh/mesh-control/control.env", + "mode": "0600", + "content": "MESH_STORE_INVENTORY=${secret:inventory}\nMESH_STORE_IDENTITY=${secret:identity}\nMESH_STORE_LICENCES=${secret:licences}\nMESH_BROKER_AMQP=${secret:broker}\nMESH_BROKER_MANAGEMENT=${secret:broker-management}\nMESH_BROKER_ADDRESS=${secret:broker-address}\n" + }, + { + "id": "server", + "type": "container", + "name": "mesh-control", + "network": "host", + "args": [ + "serve" + ], + "env-file": [ + "/var/lib/mesh/mesh-control/control.env" + ], + "env": { + "MESH_BROKER_CERTIFICATE": "/broker-tls/tls.crt" + }, + "volumes": [ + "mesh-broker-tls:/broker-tls:ro" + ], + "restart-on": [ + "control-env" + ], + "artifact": "server" + } + ], + "build": { + "artifacts": [ + { + "name": "server", + "kind": "image", + "from": "Dockerfile" + } + ] + } +}